use std::path::Path;
use std::time::{Duration, Instant};
use anyhow::{anyhow, Context as _, Result};
use colored::Colorize;
use trusty_common::launchd_claim::SocketOwner;
use trusty_mcp::StartLock;
const PROBE_TIMEOUT: Duration = Duration::from_millis(500);
const STARTUP_TIMEOUT: Duration = Duration::from_secs(30);
const POLL_INTERVAL: Duration = Duration::from_millis(500);
const LAUNCHD_RESTART_TIMEOUT: Duration =
Duration::from_secs(trusty_common::shutdown::TERMINATION_GRACE_SECS);
fn socket_owner_for(socket: &Path) -> SocketOwner {
trusty_common::launchd_claim::launchd_socket_owner(
trusty_common::launchd_labels::MEMORY,
super::single_instance::is_production_socket(socket),
)
}
async fn await_launchd_daemon(socket: &Path, label: &str, wait: Duration) -> Result<()> {
eprintln!(
"{} Waiting for launchd unit {label} to serve trusty-memory…",
"◉".cyan()
);
let deadline = Instant::now() + wait;
while Instant::now() < deadline {
tokio::time::sleep(POLL_INTERVAL).await;
if probe(socket).await {
return Ok(());
}
}
Err(anyhow!(
"launchd unit {label} owns {} but nothing served it within {}s. Not \
spawning an unsupervised daemon onto that path — it would start without \
the plist's environment and launchd's own instance would then exit 0 \
reporting success (#6619). Check `launchctl print gui/$(id -u)/{label}` \
and the daemon's launchd stderr log",
socket.display(),
wait.as_secs()
))
}
pub async fn probe(socket: &Path) -> bool {
trusty_common::uds::socket_is_serving(socket, PROBE_TIMEOUT).await
}
fn spawn_daemon() -> Result<u32> {
trusty_common::daemon_guard::spawn_current_exe_forwarding_parent_link(&[
"serve",
"--foreground",
])
.map_err(|e| anyhow!("trusty-memory daemon spawn failed: {e}"))
}
pub async fn ensure_daemon_running(socket: &Path, lock_path: &Path) -> Result<()> {
ensure_daemon_running_with(
socket,
lock_path,
socket_owner_for(socket),
LAUNCHD_RESTART_TIMEOUT,
spawn_daemon,
)
.await
}
pub(crate) async fn ensure_daemon_running_with(
socket: &Path,
lock_path: &Path,
owner: SocketOwner,
launchd_wait: Duration,
spawn: impl FnOnce() -> Result<u32>,
) -> Result<()> {
if probe(socket).await {
return Ok(());
}
let lock_path_owned = lock_path.to_path_buf();
let _lock = tokio::task::spawn_blocking(move || StartLock::acquire_blocking(&lock_path_owned))
.await
.context("start-lock acquisition task panicked")??;
if probe(socket).await {
return Ok(());
}
if let SocketOwner::Launchd { label } = &owner {
return await_launchd_daemon(socket, label, launchd_wait).await;
}
eprintln!("{} Starting trusty-memory daemon…", "◉".cyan());
spawn()?;
let deadline = Instant::now() + STARTUP_TIMEOUT;
while Instant::now() < deadline {
tokio::time::sleep(POLL_INTERVAL).await;
if probe(socket).await {
return Ok(());
}
}
Err(anyhow!(
"trusty-memory did not start serving {} within {}s — run \
`trusty-memory serve --foreground` in the foreground to see the error",
socket.display(),
STARTUP_TIMEOUT.as_secs()
))
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn probe_returns_false_for_an_absent_socket() {
let tmp = tempfile::tempdir().expect("tempdir");
let started = Instant::now();
assert!(!probe(&tmp.path().join("absent.sock")).await);
assert!(
started.elapsed() < Duration::from_secs(3),
"a refused dial must not wait out the budget: {:?}",
started.elapsed()
);
}
#[tokio::test(flavor = "multi_thread")]
async fn ensure_daemon_running_returns_early_when_something_is_serving() {
let tmp = tempfile::tempdir().expect("tempdir");
let socket = tmp.path().join("sockets").join("trusty-memory.sock");
let listener = trusty_common::uds::bind_hardened(&socket).expect("bind");
tokio::spawn(async move { while listener.accept().await.is_ok() {} });
let started = Instant::now();
ensure_daemon_running(&socket, &tmp.path().join("start.lock"))
.await
.expect("a live socket must satisfy the guard");
assert!(
started.elapsed() < Duration::from_secs(3),
"the fast path must not wait: {:?}",
started.elapsed()
);
}
#[tokio::test(flavor = "multi_thread")]
async fn ensure_daemon_running_defers_to_a_launchd_unit_instead_of_spawning() {
let tmp = tempfile::tempdir().expect("tempdir");
let err = ensure_daemon_running_with(
&tmp.path().join("absent.sock"),
&tmp.path().join("start.lock"),
SocketOwner::Launchd {
label: "com.trusty.memory".to_owned(),
},
Duration::from_millis(600),
|| unreachable!("a launchd-owned socket must never be spawned onto"),
)
.await
.expect_err("a unit that never comes back is an error, not a spawn");
assert!(
err.to_string().contains("com.trusty.memory"),
"the owning unit must be named: {err}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn ensure_daemon_running_spawns_when_no_launchd_unit_owns_the_socket() {
let tmp = tempfile::tempdir().expect("tempdir");
let socket = tmp.path().join("sockets").join("trusty-memory.sock");
let spawned = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let flag = std::sync::Arc::clone(&spawned);
let bind_at = socket.clone();
ensure_daemon_running_with(
&socket,
&tmp.path().join("start.lock"),
SocketOwner::OnDemand,
Duration::from_millis(600),
move || {
flag.store(true, std::sync::atomic::Ordering::SeqCst);
let listener = trusty_common::uds::bind_hardened(&bind_at)?;
tokio::spawn(async move { while listener.accept().await.is_ok() {} });
Ok(4242)
},
)
.await
.expect("an unmanaged host spawns and the daemon comes up");
assert!(
spawned.load(std::sync::atomic::Ordering::SeqCst),
"an unmanaged host must still get its on-demand spawn"
);
}
}