use std::collections::hash_map::DefaultHasher;
use std::fs;
use std::hash::{Hash, Hasher};
use std::io::BufRead;
use std::os::unix::net::UnixStream;
#[cfg(unix)]
use std::os::unix::process::CommandExt;
use std::path::{Path, PathBuf};
use std::process::{Child, Command, Stdio};
use std::sync::Mutex;
use std::time::{Duration, Instant};
use anyhow::{anyhow, bail, Context, Result};
#[cfg(unix)]
use libc;
use serde_json;
use super::{DaemonHandle, DaemonHost};
use crate::tui::{self, LogSource, LogTx, ServiceStatus, TuiEvent};
const HEALTH_TIMEOUT: Duration = Duration::from_secs(120);
const HEALTH_POLL: Duration = Duration::from_secs(2);
const DEV_API_PORT: u16 = 3001;
const DEV_UI_PORT: u16 = 5174;
pub struct MonorepoHost {
monorepo_path: PathBuf,
socket_override: Option<PathBuf>,
dev_dir_override: Option<PathBuf>,
log_tx: Option<LogTx>,
child: Mutex<Option<Child>>,
ui_child: Mutex<Option<Child>>,
env_file: Mutex<Option<PathBuf>>,
}
impl MonorepoHost {
pub fn new(
monorepo_path: PathBuf,
socket_override: Option<PathBuf>,
dev_dir_override: Option<PathBuf>,
log_tx: Option<LogTx>,
) -> Self {
Self {
monorepo_path,
socket_override,
dev_dir_override,
log_tx,
child: Mutex::new(None),
ui_child: Mutex::new(None),
env_file: Mutex::new(None),
}
}
fn syslog(&self, msg: impl Into<String>) {
tui::sys_log(self.log_tx.as_ref(), msg);
}
fn spawn_daemon(&self, env_file: &Path) -> Result<Child> {
let mut cmd = Command::new("cargo");
cmd.args([
"run",
"-p",
"node-server",
"--features",
"agentic_payments",
"--",
"--env",
])
.arg(env_file)
.current_dir(&self.monorepo_path)
.stdin(Stdio::null());
if self.log_tx.is_some() {
cmd.stdout(Stdio::piped()).stderr(Stdio::piped());
} else {
cmd.stdout(Stdio::inherit()).stderr(Stdio::inherit());
}
#[cfg(unix)]
cmd.process_group(0);
let mut child = cmd
.spawn()
.with_context(|| "spawn `cargo run -p node-server` — is cargo on PATH?")?;
if let Some(tx) = &self.log_tx {
if let Some(stdout) = child.stdout.take() {
let tx = tx.clone();
std::thread::spawn(move || {
for line in std::io::BufReader::new(stdout).lines().map_while(Result::ok) {
let _ = tx.send(TuiEvent::Log(crate::tui::LogEntry {
source: LogSource::Daemon,
line,
}));
}
});
}
if let Some(stderr) = child.stderr.take() {
let tx = tx.clone();
std::thread::spawn(move || {
for line in std::io::BufReader::new(stderr).lines().map_while(Result::ok) {
let _ = tx.send(TuiEvent::Log(crate::tui::LogEntry {
source: LogSource::Daemon,
line,
}));
}
});
}
}
Ok(child)
}
fn spawn_ui(&self) -> Option<Child> {
let ui_dir = self.monorepo_path.join("system/ui");
if !ui_dir.join("package.json").exists() {
return None;
}
let yarn = if which_bin("yarn") { "yarn" } else { "npm" };
self.syslog(format!(
"→ starting UI dev server ({yarn} dev --port {DEV_UI_PORT}) in {}",
ui_dir.display()
));
let mut cmd = Command::new(yarn);
cmd.args(["dev", "--port", &DEV_UI_PORT.to_string()])
.env("VITE_BACKEND_PORT", DEV_API_PORT.to_string())
.current_dir(&ui_dir)
.stdin(Stdio::null());
if self.log_tx.is_some() {
cmd.stdout(Stdio::piped()).stderr(Stdio::piped());
} else {
cmd.stdout(Stdio::inherit()).stderr(Stdio::inherit());
}
#[cfg(unix)]
cmd.process_group(0);
match cmd.spawn() {
Ok(mut child) => {
if let Some(tx) = &self.log_tx {
if let Some(stdout) = child.stdout.take() {
let tx = tx.clone();
std::thread::spawn(move || {
for line in std::io::BufReader::new(stdout).lines().map_while(Result::ok) {
let _ = tx.send(TuiEvent::Log(crate::tui::LogEntry {
source: LogSource::UiServer,
line,
}));
}
});
}
if let Some(stderr) = child.stderr.take() {
let tx = tx.clone();
std::thread::spawn(move || {
for line in std::io::BufReader::new(stderr).lines().map_while(Result::ok) {
let _ = tx.send(TuiEvent::Log(crate::tui::LogEntry {
source: LogSource::UiServer,
line,
}));
}
});
}
}
tui::update_status(
self.log_tx.as_ref(),
LogSource::UiServer,
ServiceStatus::Ready,
Some(format!("http://localhost:{DEV_UI_PORT}")),
);
self.syslog(format!("✓ UI available at http://localhost:{DEV_UI_PORT}"));
Some(child)
}
Err(e) => {
self.syslog(format!(
"⚠could not start UI dev server: {e} — run manually: \
cd system/ui && VITE_BACKEND_PORT={DEV_API_PORT} {yarn} dev --port {DEV_UI_PORT}"
));
tui::update_status(
self.log_tx.as_ref(),
LogSource::UiServer,
ServiceStatus::Disabled,
None,
);
None
}
}
}
fn shutdown_daemon_only(&self) {
if let Some(mut child) = self.child.lock().unwrap().take() {
let pid = child.id() as i32;
#[cfg(unix)]
unsafe {
libc::kill(-pid, libc::SIGTERM);
}
let t = Instant::now();
loop {
match child.try_wait() {
Ok(Some(_)) => break,
Ok(None) if t.elapsed() > Duration::from_secs(10) => {
#[cfg(unix)]
unsafe {
libc::kill(-pid, libc::SIGKILL);
}
let _ = child.wait();
break;
}
Ok(None) => std::thread::sleep(Duration::from_millis(200)),
Err(_) => break,
}
}
}
}
}
impl DaemonHost for MonorepoHost {
fn ensure_running(&self) -> Result<DaemonHandle> {
let cargo_toml = self.monorepo_path.join("Cargo.toml");
if !cargo_toml.exists() {
bail!(
"{} doesn't look like a monorepo root (no Cargo.toml found)",
self.monorepo_path.display()
);
}
let manifest = fs::read_to_string(&cargo_toml)
.with_context(|| format!("read {}", cargo_toml.display()))?;
let has_server = manifest.contains("\"system/server\"")
|| manifest.contains("system/server")
|| manifest.contains("\"apps/server\"")
|| manifest.contains("apps/server");
if !has_server {
bail!(
"{}/Cargo.toml does not include 'system/server' — is this the right path?",
self.monorepo_path.display()
);
}
let env_dir = monorepo_env_dir(&self.monorepo_path)?;
fs::create_dir_all(&env_dir).ok();
let pid_file = env_dir.join("daemon.pid");
if let Ok(contents) = fs::read_to_string(&pid_file) {
if let Ok(old_pid) = contents.trim().parse::<i32>() {
let alive = unsafe { libc::kill(old_pid, 0) == 0 };
if alive {
self.syslog(format!(
"→ previous instance found (PID {old_pid}), sending SIGTERM…"
));
unsafe {
libc::kill(-old_pid, libc::SIGTERM);
}
let deadline = Instant::now();
loop {
std::thread::sleep(Duration::from_millis(200));
if unsafe { libc::kill(old_pid, 0) } != 0 {
self.syslog(format!(
"✓ previous instance exited ({}ms)",
deadline.elapsed().as_millis()
));
break;
}
if deadline.elapsed() > Duration::from_secs(15) {
self.syslog("âš previous instance did not exit after 15s, SIGKILL");
unsafe {
libc::kill(-old_pid, libc::SIGKILL);
}
break;
}
}
}
let _ = fs::remove_file(&pid_file);
}
}
let socket_path = self
.socket_override
.clone()
.unwrap_or_else(|| env_dir.join("control.sock"));
let dev_dir = self
.dev_dir_override
.clone()
.unwrap_or_else(|| env_dir.join("dev-apps"));
let db_path = env_dir.join("dev.db");
fs::create_dir_all(&dev_dir).ok();
let env_file = env_dir.join("daemon.env");
let base_env_path = self.monorepo_path.join("system/server/alice.env");
let base_env = if base_env_path.exists() {
fs::read_to_string(&base_env_path)
.with_context(|| format!("read base env {}", base_env_path.display()))?
} else {
self.syslog(format!(
"⚠{} not found — env will be minimal",
base_env_path.display()
));
String::new()
};
let log_dir = env_dir.join("logs");
fs::create_dir_all(&log_dir).ok();
let ldk_dir = env_dir.join("ldk_data");
fs::create_dir_all(&ldk_dir).ok();
let overrides = format!(
"# --- node-app dev overrides (do not edit) ---\n\
NODE_DEV_APPS_DIR={dev_dir}\n\
DATABASE_URL=sqlite://{db}\n\
NODE_IPC_SOCKET={socket}\n\
SERVER_ADDRESS=127.0.0.1:{DEV_API_PORT}\n\
APPLICATION_LOG_DIR_PATH={log_dir}\n\
LDK_STORAGE_DIR_PATH={ldk_dir}\n\
SIGNER_SEED_PATH={ldk_dir}/signer_seed.hex\n\
ACME_ENABLED=false\n\
RUST_LOG=info,node_server::apps=debug\n\
NODE_INOTIFY_DISABLE=0\n\
STATIC_DIR_PATH=\n\n\
# --- base config (alice.env) ---\n",
dev_dir = dev_dir.display(),
db = db_path.display(),
socket = socket_path.display(),
log_dir = log_dir.display(),
ldk_dir = ldk_dir.display(),
);
fs::write(&env_file, overrides + &base_env)
.with_context(|| format!("write {}", env_file.display()))?;
*self.env_file.lock().unwrap() = Some(env_file.clone());
let _ = fs::remove_file(&socket_path);
let has_make = Command::new("make")
.arg("--version")
.stdout(Stdio::null())
.stderr(Stdio::null())
.status()
.map(|s| s.success())
.unwrap_or(false);
if has_make {
self.syslog("→ building mandatory builtin apps: make builtin-apps CARGO_PROFILE=debug");
tui::update_status(
self.log_tx.as_ref(),
LogSource::Build,
ServiceStatus::Building,
Some("builtin-apps".into()),
);
let status = Command::new("make")
.args(["builtin-apps", "CARGO_PROFILE=debug"])
.current_dir(&self.monorepo_path)
.stdin(Stdio::null())
.stdout(if self.log_tx.is_some() {
Stdio::null()
} else {
Stdio::inherit()
})
.stderr(if self.log_tx.is_some() {
Stdio::null()
} else {
Stdio::inherit()
})
.status()
.with_context(|| "spawn `make builtin-apps`")?;
if !status.success() {
self.syslog(format!(
"⚠`make builtin-apps` failed (exit {}) — some platform apps may not load",
status
));
}
} else {
self.syslog("⚠`make` not found — skipping builtin-apps build");
}
let binary_path = self.monorepo_path.join("target/debug/node-server");
if binary_path.exists() {
let _ = fs::remove_file(&binary_path);
}
self.syslog(format!(
"→ building daemon: cargo build -p node-server (cwd={})",
self.monorepo_path.display()
));
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Building,
Some(format!("cargo build (cwd={})", self.monorepo_path.display())),
);
let build_status = Command::new("cargo")
.args(["build", "-p", "node-server", "--features", "agentic_payments"])
.current_dir(&self.monorepo_path)
.stdin(Stdio::null())
.stdout(if self.log_tx.is_some() {
Stdio::null()
} else {
Stdio::inherit()
})
.stderr(if self.log_tx.is_some() {
Stdio::null()
} else {
Stdio::inherit()
})
.status()
.with_context(|| "spawn `cargo build -p node-server`")?;
if !build_status.success() {
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Failed("cargo build failed".into()),
None,
);
bail!("cargo build -p node-server failed (exit {})", build_status);
}
self.syslog(format!(
"→ spawning daemon (env={})",
env_file.display()
));
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Starting,
Some(format!("api=:{} ipc={}", DEV_API_PORT, socket_path.display())),
);
let child = self.spawn_daemon(&env_file)?;
let child_pid = child.id();
if let Err(e) = fs::write(&pid_file, child_pid.to_string()) {
self.syslog(format!("âš could not write PID file: {e}"));
}
*self.child.lock().unwrap() = Some(child);
let started = Instant::now();
loop {
if UnixStream::connect(&socket_path).is_ok() {
self.syslog(format!(
"✓ IPC socket up at {} ({}s) — waiting for apps…",
socket_path.display(),
started.elapsed().as_secs()
));
break;
}
if started.elapsed() >= HEALTH_TIMEOUT {
self.shutdown();
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Failed(format!(
"socket not connectable after {}s",
HEALTH_TIMEOUT.as_secs()
)),
None,
);
return Err(anyhow!(
"daemon did not come up within {}s (socket at {} not connectable).",
HEALTH_TIMEOUT.as_secs(),
socket_path.display()
));
}
if let Some(child) = self.child.lock().unwrap().as_mut() {
if let Ok(Some(status)) = child.try_wait() {
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Failed(format!("exited {}", status)),
None,
);
return Err(anyhow!("daemon exited with {} before becoming ready", status));
}
}
std::thread::sleep(HEALTH_POLL);
}
*self.ui_child.lock().unwrap() = self.spawn_ui();
let apps_started = Instant::now();
self.syslog("→ waiting for critical apps (device-registry, core-storage)…");
loop {
let all_ready = ipc_list_app_statuses(&socket_path)
.map(|statuses| {
["device-registry", "core-storage"].iter().all(|name| {
statuses
.get(*name)
.map(|s| s == "active" || s == "running" || s == "lazy")
.unwrap_or(false)
})
})
.unwrap_or(false);
if all_ready {
self.syslog(format!(
"✓ critical apps ready ({}s)",
apps_started.elapsed().as_secs()
));
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Ready,
Some(format!(
"api=http://localhost:{DEV_API_PORT} ui=http://localhost:{DEV_UI_PORT}"
)),
);
break;
}
if apps_started.elapsed() >= Duration::from_secs(60) {
self.syslog("⚠critical apps not fully loaded after 60s — proceeding anyway");
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Ready,
Some("api ready (some apps slow)".into()),
);
break;
}
if let Some(child) = self.child.lock().unwrap().as_mut() {
if let Ok(Some(status)) = child.try_wait() {
return Err(anyhow!("daemon exited with {} while waiting for apps", status));
}
}
std::thread::sleep(Duration::from_secs(2));
}
Ok(DaemonHandle {
banner: format!(
"monorepo daemon (cargo run from {}, api=http://localhost:{DEV_API_PORT}, \
ui=http://localhost:{DEV_UI_PORT}, socket={})",
self.monorepo_path.display(),
socket_path.display()
),
socket_path,
dev_dir,
})
}
fn tail_logs(&self, app_name: &str) {
let log_tx = match &self.log_tx {
Some(tx) => tx.clone(),
None => return,
};
let env_dir = match monorepo_env_dir(&self.monorepo_path) {
Ok(d) => d,
Err(_) => return,
};
let log_file = env_dir
.join("ldk_data")
.join("logs")
.join("apps")
.join(app_name)
.join("app.log");
std::thread::spawn(move || {
use std::io::{BufRead, BufReader, Seek, SeekFrom};
let deadline = std::time::Instant::now();
loop {
if log_file.exists() {
break;
}
if deadline.elapsed() > std::time::Duration::from_secs(5) {
return;
}
std::thread::sleep(std::time::Duration::from_millis(200));
}
let file = match fs::File::open(&log_file) {
Ok(f) => f,
Err(_) => return,
};
let mut reader = BufReader::new(file);
let _ = reader.seek(SeekFrom::Start(0));
loop {
let mut line = String::new();
match reader.read_line(&mut line) {
Ok(0) => {
std::thread::sleep(std::time::Duration::from_millis(100));
}
Ok(_) => {
let trimmed =
line.trim_end_matches('\n').trim_end_matches('\r').to_string();
if !trimmed.is_empty()
&& log_tx
.send(TuiEvent::Log(crate::tui::LogEntry {
source: LogSource::App,
line: trimmed,
}))
.is_err()
{
return; }
}
Err(_) => return,
}
}
});
}
fn shutdown(&self) {
if let Ok(env_dir) = monorepo_env_dir(&self.monorepo_path) {
let _ = fs::remove_file(env_dir.join("daemon.pid"));
}
if let Some(mut child) = self.child.lock().unwrap().take() {
let pid = child.id() as i32;
self.syslog(format!("→ shutting down monorepo daemon (PID {pid})…"));
#[cfg(unix)]
unsafe {
libc::kill(-pid, libc::SIGTERM);
}
let t = Instant::now();
loop {
match child.try_wait() {
Ok(Some(_)) => {
self.syslog(format!(
"✓ daemon exited ({}ms)",
t.elapsed().as_millis()
));
break;
}
Ok(None) if t.elapsed() >= Duration::from_secs(10) => {
self.syslog("âš daemon did not exit after 10s, SIGKILL");
#[cfg(unix)]
unsafe {
libc::kill(-pid, libc::SIGKILL);
}
let _ = child.wait();
break;
}
Ok(None) => std::thread::sleep(Duration::from_millis(200)),
Err(e) => {
self.syslog(format!("âš daemon wait error: {e}"));
break;
}
}
}
}
if let Some(mut ui) = self.ui_child.lock().unwrap().take() {
let pid = ui.id() as i32;
self.syslog(format!("→ shutting down UI dev server (PID {pid})…"));
#[cfg(unix)]
unsafe {
libc::kill(-pid, libc::SIGTERM);
}
let t = Instant::now();
loop {
match ui.try_wait() {
Ok(Some(_)) => break,
Ok(None) if t.elapsed() >= Duration::from_secs(5) => {
#[cfg(unix)]
unsafe {
libc::kill(-pid, libc::SIGKILL);
}
let _ = ui.wait();
break;
}
Ok(None) => std::thread::sleep(Duration::from_millis(200)),
Err(_) => break,
}
}
}
}
fn restart(&self) -> Result<()> {
let env_file = self
.env_file
.lock()
.unwrap()
.clone()
.ok_or_else(|| anyhow!("daemon was never started — cannot restart"))?;
self.syslog("→ restart: stopping daemon…");
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Starting,
Some("restarting…".into()),
);
self.shutdown_daemon_only();
let socket_path = self
.socket_override
.clone()
.unwrap_or_else(|| {
monorepo_env_dir(&self.monorepo_path)
.ok()
.map(|d| d.join("control.sock"))
.unwrap_or_else(|| PathBuf::from("/tmp/node-control.sock"))
});
let _ = fs::remove_file(&socket_path);
self.syslog("→ restart: spawning daemon…");
let child = self.spawn_daemon(&env_file)?;
let child_pid = child.id();
*self.child.lock().unwrap() = Some(child);
if let Ok(env_dir) = monorepo_env_dir(&self.monorepo_path) {
let _ = fs::write(env_dir.join("daemon.pid"), child_pid.to_string());
}
let t = Instant::now();
loop {
if UnixStream::connect(&socket_path).is_ok() {
self.syslog(format!(
"✓ daemon restarted ({}s)",
t.elapsed().as_secs()
));
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Ready,
Some("restarted".into()),
);
return Ok(());
}
if t.elapsed() > HEALTH_TIMEOUT {
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Failed("socket not connectable after restart".into()),
None,
);
bail!("daemon did not come up after restart");
}
if let Some(child) = self.child.lock().unwrap().as_mut() {
if let Ok(Some(status)) = child.try_wait() {
tui::update_status(
self.log_tx.as_ref(),
LogSource::Daemon,
ServiceStatus::Failed(format!("exited {}", status)),
None,
);
bail!("daemon exited with {} during restart", status);
}
}
std::thread::sleep(HEALTH_POLL);
}
}
fn pre_start_dev_dir(&self) -> Option<PathBuf> {
monorepo_dev_dir(&self.monorepo_path, self.dev_dir_override.as_deref()).ok()
}
}
fn monorepo_env_dir(monorepo_path: &Path) -> Result<PathBuf> {
let cache_dir = cache_root()?;
let path_hash = {
let mut h = DefaultHasher::new();
monorepo_path
.canonicalize()
.unwrap_or_else(|_| monorepo_path.to_path_buf())
.hash(&mut h);
h.finish()
};
Ok(cache_dir.join(format!("monorepo-{:x}", path_hash)))
}
fn monorepo_dev_dir(monorepo_path: &Path, dev_dir_override: Option<&Path>) -> Result<PathBuf> {
if let Some(override_path) = dev_dir_override {
return Ok(override_path.to_path_buf());
}
Ok(monorepo_env_dir(monorepo_path)?.join("dev-apps"))
}
fn ipc_list_app_statuses(
socket_path: &std::path::Path,
) -> Result<std::collections::HashMap<String, String>> {
use std::io::{BufRead, BufReader, Write};
let mut stream = UnixStream::connect(socket_path)
.with_context(|| "connect to IPC socket")?;
stream.set_read_timeout(Some(Duration::from_secs(5))).ok();
stream.set_write_timeout(Some(Duration::from_secs(5))).ok();
let request = serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "app.list",
"params": {}
});
let mut line = serde_json::to_string(&request).unwrap();
line.push('\n');
stream.write_all(line.as_bytes())?;
let reader = BufReader::new(&stream);
let response_line = reader
.lines()
.next()
.ok_or_else(|| anyhow!("no response"))??;
let v: serde_json::Value = serde_json::from_str(&response_line)?;
let mut map = std::collections::HashMap::new();
if let Some(apps) = v.pointer("/result/apps").and_then(|a| a.as_array()) {
for app in apps {
if let (Some(name), Some(status)) = (
app.get("name").and_then(|n| n.as_str()),
app.get("status").and_then(|s| s.as_str()),
) {
map.insert(name.to_string(), status.to_string());
}
}
}
Ok(map)
}
fn which_bin(bin: &str) -> bool {
std::env::var_os("PATH")
.map(|path| std::env::split_paths(&path).any(|dir| dir.join(bin).is_file()))
.unwrap_or(false)
}
fn cache_root() -> Result<PathBuf> {
if let Ok(c) = std::env::var("XDG_CACHE_HOME") {
if !c.is_empty() {
return Ok(PathBuf::from(c).join("node-app"));
}
}
let home = std::env::var_os("HOME").ok_or_else(|| anyhow!("$HOME not set"))?;
Ok(PathBuf::from(home).join(".cache/node-app"))
}