use std::ffi::OsString;
use std::path::Path;
use std::time::{Duration, Instant};
use mcpmesh::client::connect_control;
use serde_json::json;
pub const STUB: &str = env!("CARGO_BIN_EXE_echo_mcp_stub");
pub const MCPMESH: &str = env!("CARGO_BIN_EXE_mcpmesh");
pub fn world(dir: &Path, config_body: &str) -> (std::path::PathBuf, Vec<(OsString, OsString)>) {
let runtime = dir.join("runtime");
let config = dir.join("config");
let data = dir.join("data");
let config_mcpmesh = config.join("mcpmesh");
std::fs::create_dir_all(&config_mcpmesh).unwrap();
std::fs::write(config_mcpmesh.join("config.toml"), config_body).unwrap();
let socket = runtime.join("mcpmesh").join("mcpmesh.sock");
let env = vec![
(OsString::from("HOME"), dir.as_os_str().to_os_string()),
(OsString::from("XDG_RUNTIME_DIR"), runtime.into_os_string()),
(OsString::from("XDG_CONFIG_HOME"), config.into_os_string()),
(OsString::from("XDG_DATA_HOME"), data.into_os_string()),
];
(socket, env)
}
pub fn run_cmd(env: &[(OsString, OsString)], args: &[&str]) -> std::process::Output {
let mut cmd = std::process::Command::new(MCPMESH);
cmd.args(args);
for (k, v) in env {
cmd.env(k, v);
}
cmd.output().expect("run mcpmesh subcommand")
}
pub async fn shutdown_daemon(socket: &Path) {
if let Ok(mut client) = connect_control(socket).await {
let _ = client.request_value(&json!({ "method": "shutdown" })).await;
}
let deadline = Instant::now() + Duration::from_secs(10);
loop {
if connect_control(socket).await.is_err() && singleton_lock_is_free(socket) {
return;
}
assert!(
Instant::now() < deadline,
"daemon still accepting connections (or holding the singleton flock) after shutdown"
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
fn singleton_lock_is_free(socket: &Path) -> bool {
use rustix::fs::{FlockOperation, flock};
let lock_path = socket.with_file_name("mcpmesh.lock");
let Ok(file) = std::fs::OpenOptions::new().read(true).open(&lock_path) else {
return !lock_path.exists();
};
flock(&file, FlockOperation::NonBlockingLockExclusive).is_ok()
}