use crate::drift::{classify, DriftState, ExeFingerprint};
use crate::paths::AgentsHome;
use crate::protocol::{read_response, write_request, ProtocolError, Request, Response};
use serde_json::{json, Value};
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use tokio::net::UnixStream;
#[derive(Debug, thiserror::Error)]
pub enum ClientError {
#[error("protocol: {0}")]
Protocol(#[from] ProtocolError),
#[error("io: {0}")]
Io(#[from] std::io::Error),
#[error("daemon did not come up within {0:?}")]
DaemonStartTimeout(Duration),
#[error(
"daemon binary not found: {0} - the fno-agents triad (client/daemon/worker) \
is split here. Run `fno update` to redeploy the pair, or set \
FNO_AGENTS_DAEMON_BIN to a coherent same-build daemon."
)]
DaemonBinMissing(PathBuf),
#[error("daemon is not running")]
DaemonNotRunning,
}
pub fn resolve_daemon_bin() -> PathBuf {
if let Some(v) = std::env::var_os("FNO_AGENTS_DAEMON_BIN") {
return PathBuf::from(v);
}
std::env::current_exe()
.ok()
.and_then(|p| p.parent().map(|d| d.join("fno-agents-daemon")))
.unwrap_or_else(|| PathBuf::from("fno-agents-daemon"))
}
pub async fn ensure_daemon(
home: &AgentsHome,
daemon_bin: &std::path::Path,
) -> Result<(), ClientError> {
let sock = home.supervisor_sock();
if UnixStream::connect(&sock).await.is_ok() {
return Ok(());
}
if !daemon_bin.exists() {
return Err(ClientError::DaemonBinMissing(daemon_bin.to_path_buf()));
}
eprintln!("(lazy-starting daemon)");
let mut cmd = tokio::process::Command::new(daemon_bin);
cmd.process_group(0);
cmd.env("FNO_AGENTS_HOME", home.root());
if let Some(v) = std::env::var_os("FNO_AGENTS_WORKER_BIN") {
cmd.env("FNO_AGENTS_WORKER_BIN", v);
}
let child = cmd.spawn()?;
drop(child);
let start = Instant::now();
let budget = Duration::from_secs(10);
while start.elapsed() < budget {
if UnixStream::connect(&sock).await.is_ok() {
return Ok(());
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
Err(ClientError::DaemonStartTimeout(budget))
}
pub async fn call(
home: &AgentsHome,
daemon_bin: &std::path::Path,
req: &Request,
) -> Result<Response, ClientError> {
ensure_daemon(home, daemon_bin).await?;
let mut conn = UnixStream::connect(home.supervisor_sock()).await?;
write_request(&mut conn, req).await?;
Ok(read_response(&mut conn).await?)
}
pub async fn call_if_running(home: &AgentsHome, req: &Request) -> Result<Response, ClientError> {
let mut conn = match UnixStream::connect(home.supervisor_sock()).await {
Ok(c) => c,
Err(e)
if matches!(
e.kind(),
std::io::ErrorKind::NotFound | std::io::ErrorKind::ConnectionRefused
) =>
{
return Err(ClientError::DaemonNotRunning)
}
Err(e) => return Err(ClientError::Io(e)),
};
write_request(&mut conn, req).await?;
Ok(read_response(&mut conn).await?)
}
fn running_fingerprint(status: &Value) -> Option<ExeFingerprint> {
let d = status.get("daemon")?;
let path = d.get("exe_path")?.as_str()?;
let mtime_nanos = d.get("exe_mtime")?.as_i64()?;
let size = d.get("exe_size")?.as_u64()?;
Some(ExeFingerprint {
path: PathBuf::from(path),
mtime_nanos,
size,
})
}
pub fn drift_from_status(status: &Value) -> DriftState {
let running = running_fingerprint(status);
let on_disk = ExeFingerprint::of(&resolve_daemon_bin());
classify(running.as_ref(), on_disk.as_ref())
}
pub async fn check_daemon_drift(home: &AgentsHome) -> DriftState {
let req = Request::new(1, "agent.status", json!({}));
match call_if_running(home, &req).await {
Ok(resp) => match resp.result() {
Some(result) => drift_from_status(result),
None => DriftState::Unknown,
},
Err(ClientError::DaemonNotRunning) => DriftState::DaemonDown,
Err(_) => DriftState::Unknown,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RestartOutcome {
pub old_pid: Option<u32>,
pub new_pid: u32,
}
#[derive(Debug, thiserror::Error)]
pub enum RestartError {
#[error("SIGTERM to daemon pid {pid} failed: {reason}")]
SigtermFailed { pid: u32, reason: String },
#[error("daemon pid {pid} did not exit after SIGTERM within {secs}s; check it manually")]
DidNotExit { pid: u32, secs: u64 },
#[error("daemon status response missing daemon.pid")]
StatusMissingPid,
#[error(transparent)]
Client(#[from] ClientError),
}
const RESTART_SOCKET_TIMEOUT: Duration = Duration::from_secs(5);
enum SigtermResult {
Sent,
AlreadyGone,
Failed(String),
}
fn send_sigterm(pid: u32) -> SigtermResult {
let rc = unsafe { libc::kill(pid as libc::pid_t, libc::SIGTERM) };
if rc == 0 {
return SigtermResult::Sent;
}
let err = std::io::Error::last_os_error();
match err.raw_os_error() {
Some(e) if e == libc::ESRCH => SigtermResult::AlreadyGone,
_ => SigtermResult::Failed(err.to_string()),
}
}
async fn read_daemon_pid(home: &AgentsHome) -> Result<u32, RestartError> {
let req = Request::new(1, "agent.status", json!({}));
let resp = call_if_running(home, &req).await?;
resp.result()
.and_then(|r| r.get("daemon"))
.and_then(|d| d.get("pid"))
.and_then(Value::as_u64)
.map(|p| p as u32)
.ok_or(RestartError::StatusMissingPid)
}
async fn await_socket_clear(home: &AgentsHome) -> bool {
let sock = home.supervisor_sock();
let start = Instant::now();
while start.elapsed() < RESTART_SOCKET_TIMEOUT {
if UnixStream::connect(&sock).await.is_err() {
return true;
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
UnixStream::connect(&sock).await.is_err()
}
async fn start_fresh(home: &AgentsHome, daemon_bin: &Path) -> Result<u32, RestartError> {
ensure_daemon(home, daemon_bin).await?;
read_daemon_pid(home).await
}
pub async fn restart_daemon(
home: &AgentsHome,
daemon_bin: &Path,
) -> Result<RestartOutcome, RestartError> {
let status = Request::new(1, "agent.status", json!({}));
let result = match call_if_running(home, &status).await {
Ok(resp) => resp.result().cloned(),
Err(ClientError::DaemonNotRunning) => None,
Err(e) => return Err(RestartError::Client(e)),
};
let Some(result) = result else {
let new_pid = start_fresh(home, daemon_bin).await?;
return Ok(RestartOutcome {
old_pid: None,
new_pid,
});
};
let daemon = result.get("daemon");
let old_pid = daemon
.and_then(|d| d.get("pid"))
.and_then(Value::as_u64)
.ok_or(RestartError::StatusMissingPid)? as u32;
let recorded_start = daemon
.and_then(|d| d.get("pid_start_time"))
.and_then(Value::as_u64);
if !crate::daemon::pid_is_ours(old_pid, recorded_start) {
let new_pid = start_fresh(home, daemon_bin).await?;
return Ok(RestartOutcome {
old_pid: Some(old_pid),
new_pid,
});
}
match send_sigterm(old_pid) {
SigtermResult::Sent | SigtermResult::AlreadyGone => {}
SigtermResult::Failed(reason) => {
return Err(RestartError::SigtermFailed {
pid: old_pid,
reason,
})
}
}
if !await_socket_clear(home).await {
return Err(RestartError::DidNotExit {
pid: old_pid,
secs: RESTART_SOCKET_TIMEOUT.as_secs(),
});
}
let new_pid = start_fresh(home, daemon_bin).await?;
Ok(RestartOutcome {
old_pid: Some(old_pid),
new_pid,
})
}
#[cfg(test)]
mod drift_restart_tests {
use super::*;
#[test]
fn drift_from_status_unknown_when_no_fingerprint() {
let status = json!({"daemon": {"pid": 1, "state": "serving"}});
assert_eq!(drift_from_status(&status), DriftState::Unknown);
}
#[test]
fn running_fingerprint_parses_and_rejects() {
let status = json!({"daemon": {
"exe_path": "/opt/fno-agents-daemon",
"exe_mtime": 1_700_000_000_000_000_000_i64,
"exe_size": 4242_u64,
}});
let fp = running_fingerprint(&status).expect("parses a full fingerprint");
assert_eq!(fp.path, PathBuf::from("/opt/fno-agents-daemon"));
assert_eq!(fp.mtime_nanos, 1_700_000_000_000_000_000);
assert_eq!(fp.size, 4242);
let partial = json!({"daemon": {"exe_path": "/opt/x", "exe_size": 1_u64}});
assert!(running_fingerprint(&partial).is_none());
assert!(running_fingerprint(&json!({})).is_none());
}
#[test]
fn send_sigterm_reports_already_gone_for_dead_pid() {
let child = std::process::Command::new("true")
.spawn()
.expect("spawn true");
let pid = child.id();
let mut child = child;
let _ = child.wait(); std::thread::sleep(Duration::from_millis(50));
match send_sigterm(pid) {
SigtermResult::AlreadyGone | SigtermResult::Sent => {}
SigtermResult::Failed(e) => panic!("unexpected Failed: {e}"),
}
}
}