use std::ffi::OsString;
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use anyhow::{Context, Result};
use mcpmesh_trust::paths;
pub use mcpmesh_local_api::client::{ClientError, ControlClient, connect_control};
#[derive(Clone, Debug)]
pub struct DaemonLaunch {
pub exe: PathBuf,
pub socket: PathBuf,
pub env: Vec<(OsString, OsString)>,
}
impl DaemonLaunch {
pub fn ambient() -> Result<Self> {
let mut env = Vec::new();
if let Some(root) = paths::profile_root() {
env.push((
std::ffi::OsString::from("MCPMESH_HOME"),
root.into_os_string(),
));
}
Ok(Self {
exe: std::env::current_exe().context("resolve current executable")?,
socket: paths::default_endpoint()?,
env,
})
}
}
pub async fn ensure_daemon() -> Result<ControlClient> {
ensure_daemon_with(&DaemonLaunch::ambient()?).await
}
pub const DEFAULT_READY_TIMEOUT: Duration = Duration::from_secs(10);
pub async fn ensure_daemon_with(launch: &DaemonLaunch) -> Result<ControlClient> {
ensure_daemon_with_timeout(launch, DEFAULT_READY_TIMEOUT).await
}
pub async fn ensure_daemon_with_timeout(
launch: &DaemonLaunch,
ready_timeout: Duration,
) -> Result<ControlClient> {
#[cfg(unix)]
check_socket_path_len(&launch.socket)?;
if let Ok(client) = connect_control(&launch.socket).await {
return Ok(client);
}
let mut spawned = spawn_detached(launch)?;
let deadline = Instant::now() + ready_timeout;
let mut backoff = Duration::from_millis(20);
loop {
match connect_control(&launch.socket).await {
Ok(client) => return Ok(client),
Err(e) => {
if let Ok(Some(status)) = spawned.child.try_wait()
&& !status.success()
{
anyhow::bail!("{}", autostart_failure_message(&spawned.read_stderr()));
}
if Instant::now() >= deadline {
let stderr = spawned.read_stderr();
if !stderr.is_empty() {
anyhow::bail!("{}", autostart_failure_message(&stderr));
}
return Err(anyhow::Error::from(e).context(format!(
"the daemon did not accept connections within {}s — run \
'mcpmesh doctor' to diagnose",
ready_timeout.as_secs()
)));
}
tokio::time::sleep(backoff).await;
backoff = (backoff * 2).min(Duration::from_millis(200));
}
}
}
}
#[cfg(any(target_os = "linux", target_os = "android"))]
const MAX_SOCKET_PATH: usize = 107;
#[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
const MAX_SOCKET_PATH: usize = 103;
#[cfg(unix)]
fn check_socket_path_len(socket: &Path) -> Result<()> {
use std::os::unix::ffi::OsStrExt;
if socket.as_os_str().as_bytes().len() > MAX_SOCKET_PATH {
anyhow::bail!(
"runtime dir path too long for a unix socket — set XDG_RUNTIME_DIR to a shorter path"
);
}
Ok(())
}
fn autostart_failure_message(stderr: &str) -> String {
let reason = stderr.trim();
let reason = reason.strip_prefix("Error: ").unwrap_or(reason);
if reason.is_empty() {
"the daemon failed to start — run 'mcpmesh doctor' to diagnose".to_string()
} else {
format!("the daemon failed to start: {reason}\nrun 'mcpmesh doctor' to diagnose")
}
}
struct SpawnedDaemon {
child: std::process::Child,
stderr_path: Option<PathBuf>,
}
impl SpawnedDaemon {
fn read_stderr(&self) -> String {
self.stderr_path
.as_deref()
.and_then(|p| std::fs::read_to_string(p).ok())
.map(|s| s.trim().to_string())
.unwrap_or_default()
}
}
impl Drop for SpawnedDaemon {
fn drop(&mut self) {
if let Some(p) = &self.stderr_path {
let _ = std::fs::remove_file(p);
}
}
}
fn stderr_capture(socket: &Path) -> Option<(std::fs::File, PathBuf)> {
use std::sync::atomic::{AtomicU64, Ordering};
static SEQ: AtomicU64 = AtomicU64::new(0);
let name = format!(
"daemon-start-{}-{}.stderr",
std::process::id(),
SEQ.fetch_add(1, Ordering::Relaxed)
);
let path = capture_dir(socket).join(name);
let file = std::fs::File::create(&path).ok()?;
Some((file, path))
}
#[cfg(unix)]
fn capture_dir(socket: &Path) -> PathBuf {
match socket.parent() {
Some(parent) if crate::ipc::ensure_runtime_dir(parent).is_ok() => parent.to_path_buf(),
_ => std::env::temp_dir(),
}
}
#[cfg(windows)]
fn capture_dir(_socket: &Path) -> PathBuf {
std::env::temp_dir()
}
fn spawn_detached(launch: &DaemonLaunch) -> Result<SpawnedDaemon> {
let (stderr_stdio, stderr_path) = match stderr_capture(&launch.socket) {
Some((file, path)) => (std::process::Stdio::from(file), Some(path)),
None => (std::process::Stdio::null(), None),
};
let mut cmd = std::process::Command::new(&launch.exe);
cmd.arg("internal")
.arg("daemon")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(stderr_stdio);
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
cmd.process_group(0);
}
#[cfg(windows)]
{
use std::os::windows::process::CommandExt;
cmd.creation_flags(0x0000_0008 | 0x0000_0200);
}
for (key, value) in &launch.env {
cmd.env(key, value);
}
match cmd.spawn() {
Ok(child) => Ok(SpawnedDaemon { child, stderr_path }),
Err(e) => {
if let Some(p) = &stderr_path {
let _ = std::fs::remove_file(p);
}
Err(anyhow::Error::from(e)
.context(format!("spawn {} internal daemon", launch.exe.display())))
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn autostart_failure_message_replays_the_daemons_reason_and_names_doctor() {
let msg = autostart_failure_message(
"Error: config error in /home/x/config.toml: [network] relay_mode = \"custom\" \
requires at least one relay_urls entry\n",
);
assert!(
msg.starts_with("the daemon failed to start: config error"),
"the daemon's own reason leads: {msg}"
);
assert!(
!msg.contains("Error:"),
"no stuttered Error: prefix inside the message: {msg}"
);
assert!(
msg.contains("run 'mcpmesh doctor' to diagnose"),
"the exact next command is named: {msg}"
);
}
#[test]
fn autostart_failure_message_degrades_cleanly_with_nothing_captured() {
let msg = autostart_failure_message("");
assert_eq!(
msg,
"the daemon failed to start — run 'mcpmesh doctor' to diagnose"
);
}
#[cfg(unix)]
#[tokio::test]
async fn overlong_socket_path_is_refused_immediately_in_user_language() {
let long_dir = std::env::temp_dir().join("x".repeat(200));
let launch = DaemonLaunch {
exe: PathBuf::from("/nonexistent-mcpmesh"),
socket: long_dir.join("mcpmesh").join("mcpmesh.sock"),
env: Vec::new(),
};
let start = Instant::now();
let err = ensure_daemon_with(&launch).await.expect_err("must refuse");
let msg = err.to_string();
assert!(
msg.contains("runtime dir path too long for a unix socket")
&& msg.contains("set XDG_RUNTIME_DIR to a shorter path"),
"the refusal names the cause and the fix: {msg}"
);
assert!(
!msg.contains("SUN_LEN"),
"no kernel vocabulary reaches the user: {msg}"
);
assert!(
start.elapsed() < Duration::from_secs(2),
"the refusal must not wait out the connect window"
);
}
#[cfg(unix)]
#[test]
fn short_socket_path_passes_the_len_check() {
assert!(check_socket_path_len(Path::new("/run/user/501/mcpmesh/mcpmesh.sock")).is_ok());
}
}