use super::Supervisor;
use super::hooks::{HookType, fire_hook};
use crate::daemon::Daemon;
use crate::daemon_id::DaemonId;
use crate::daemon_status::DaemonStatus;
use crate::procs::PROCS;
use crate::settings::settings;
use crate::supervisor::SUPERVISOR;
use crate::supervisor::state::UpsertDaemonOpts;
use std::sync::atomic;
use std::time::Duration;
use tokio::time;
const POLL_INTERVAL: Duration = Duration::from_secs(2);
enum PollOutcome {
ProcessDied { was_stopping: bool },
StoppedExternally,
TakenOver,
}
static NEXT_MONITOR_TOKEN: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) struct MonitorEntry {
pub(crate) pid: u32,
pub(crate) token: u64,
}
pub(crate) struct MonitoredGuard {
id: DaemonId,
token: u64,
}
impl MonitoredGuard {
pub(crate) fn register(id: DaemonId, pid: u32) -> Self {
let token = NEXT_MONITOR_TOKEN.fetch_add(1, atomic::Ordering::Relaxed);
SUPERVISOR
.monitored
.lock()
.expect("monitored lock poisoned")
.insert(id.clone(), MonitorEntry { pid, token });
Self { id, token }
}
pub(crate) fn try_register(id: DaemonId, pid: u32) -> Option<Self> {
let mut monitored = SUPERVISOR
.monitored
.lock()
.expect("monitored lock poisoned");
if monitored.contains_key(&id) {
return None;
}
let token = NEXT_MONITOR_TOKEN.fetch_add(1, atomic::Ordering::Relaxed);
monitored.insert(id.clone(), MonitorEntry { pid, token });
Some(Self { id, token })
}
pub(crate) fn token(&self) -> u64 {
self.token
}
}
impl Drop for MonitoredGuard {
fn drop(&mut self) {
let mut monitored = SUPERVISOR
.monitored
.lock()
.expect("monitored lock poisoned");
if monitored.get(&self.id).map(|e| e.token) == Some(self.token) {
monitored.remove(&self.id);
}
}
}
impl Supervisor {
pub(crate) fn is_monitored(&self, id: &DaemonId, pid: u32) -> bool {
self.monitored
.lock()
.expect("monitored lock poisoned")
.get(id)
.is_some_and(|e| e.pid == pid)
}
fn monitor_token_valid(&self, id: &DaemonId, token: u64) -> bool {
self.monitored
.lock()
.expect("monitored lock poisoned")
.get(id)
.is_some_and(|e| e.token == token)
}
pub(crate) async fn reconcile_unmonitored_daemons(&self) {
if !settings().supervisor.cleanup_orphans {
return;
}
let candidates: Vec<Daemon> = {
let state = self.state_file.lock().await;
state
.daemons
.values()
.filter(|d| {
d.id != DaemonId::pitchfork()
&& d.status.is_running()
&& d.pid.is_some_and(|pid| !self.is_monitored(&d.id, pid))
})
.cloned()
.collect()
};
if candidates.is_empty() {
return;
}
let policy = super::orphan_policy();
let boot_time = PROCS.boot_time();
for daemon in candidates {
let Some(pid) = daemon.pid else { continue };
if self.is_monitored(&daemon.id, pid) {
continue;
}
PROCS.refresh_pids(&[pid]);
if !PROCS.is_running(pid) {
let status = super::unobserved_exit_status(
&daemon.status,
daemon.boot_time,
boot_time,
true,
);
warn!(
"daemon {} (pid {pid}) died while unmonitored; marking {status}",
daemon.id
);
self.finalize_if_pid(&daemon.id, pid, status, super::ExitObservation::Unobserved)
.await;
continue;
}
let current_start_time = PROCS.start_time(pid);
let current_title = PROCS.title(pid);
let matches = super::process_identity_matches(
daemon.start_time,
daemon.title.as_deref(),
current_start_time,
current_title.as_deref(),
);
if !matches {
if daemon.start_time.is_some() && current_start_time.is_none() {
warn!(
"could not verify start time for live pid {pid} recorded for daemon {}; retaining running state",
daemon.id,
);
continue;
}
let status = super::unobserved_exit_status(
&daemon.status,
daemon.boot_time,
boot_time,
true,
);
warn!(
"pid {pid} recorded for daemon {} belongs to a different process now (PID recycled); marking {status}",
daemon.id,
);
self.finalize_if_pid(&daemon.id, pid, status, super::ExitObservation::Unobserved)
.await;
continue;
}
let Some(expected_start_time) = current_start_time else {
warn!(
"could not read start time for live pid {pid} recorded for daemon {}; retaining running state",
daemon.id,
);
continue;
};
if policy == "adopt" {
self.adopt_daemon(&daemon, pid, expected_start_time).await;
continue;
}
info!(
"terminating unmonitored orphaned daemon {} (pid {pid})",
daemon.id
);
let stop_cfg = daemon.stop_signal.unwrap_or_default();
match PROCS
.kill_process_group_if_start_time_matches_async(
pid,
expected_start_time,
stop_cfg.signal.into(),
stop_cfg.timeout,
)
.await
{
Ok(true) => {
self.finalize_if_pid(
&daemon.id,
pid,
DaemonStatus::Stopped,
super::ExitObservation::Terminated,
)
.await;
}
Ok(false) => {
warn!(
"could not securely terminate unmonitored orphaned daemon {} (pid {pid}); retaining running state",
daemon.id
);
}
Err(err) => {
warn!(
"failed to terminate unmonitored orphaned daemon {} (pid {pid}): {err}; retaining running state",
daemon.id
);
}
}
}
}
pub(crate) async fn finalize_monitored_exit(
&self,
id: &DaemonId,
pid: u32,
token: u64,
status: DaemonStatus,
last_exit_success: Option<bool>,
) -> bool {
let mut state_file = self.state_file.lock().await;
if !self.monitor_token_valid(id, token) {
debug!("daemon {id} was superseded; skipping exit finalization");
return false;
}
let Some(d) = state_file.daemons.get(id) else {
return false;
};
if d.pid != Some(pid) {
debug!("daemon {id} was claimed by a successor; skipping exit finalization");
return false;
}
let mut d = d.clone();
d.pid = None;
d.title = None;
d.start_time = None;
d.boot_time = None;
d.status = status;
d.last_exit_success = last_exit_success;
d.active_port = None;
state_file.clear_active_port(id);
state_file.insert_daemon(id, d);
true
}
async fn finalize_if_pid(
&self,
id: &DaemonId,
pid: u32,
status: DaemonStatus,
observation: super::ExitObservation,
) -> bool {
let mut state_file = self.state_file.lock().await;
if self
.monitored
.lock()
.expect("monitored lock poisoned")
.contains_key(id)
{
debug!("daemon {id} gained a monitor since the snapshot; skipping finalization");
return false;
}
let Some(d) = state_file.daemons.get(id) else {
return false;
};
if d.pid != Some(pid) {
debug!("daemon {id} was claimed by a successor; skipping finalization");
return false;
}
let mut d = d.clone();
d.pid = None;
d.title = None;
d.start_time = None;
d.boot_time = None;
d.status = status;
d.last_exit_success = observation.last_exit_success();
d.active_port = None;
state_file.clear_active_port(id);
state_file.insert_daemon(id, d);
true
}
pub(crate) async fn adopt_daemon(&self, daemon: &Daemon, pid: u32, expected_start_time: u64) {
let Some(guard) = MonitoredGuard::try_register(daemon.id.clone(), pid) else {
debug!(
"daemon {} (pid {pid}) is already monitored; skipping duplicate adoption",
daemon.id
);
return;
};
info!("re-adopting orphaned daemon {} (pid {pid})", daemon.id);
if daemon.start_time.is_none() {
let active_port = daemon.active_port;
let _ = self
.upsert_daemon(
UpsertDaemonOpts::builder(daemon.id.clone())
.set(|o| {
o.pid = Some(pid);
o.status = DaemonStatus::Running;
o.active_port = active_port;
})
.build(),
)
.await;
}
let id = daemon.id.clone();
let daemon_dir = daemon
.dir
.clone()
.unwrap_or_else(|| crate::env::CWD.clone());
let hook_env = daemon.env.clone();
let hook_retry = daemon.retry;
let hook_retry_count = daemon.retry_count;
tokio::spawn(async move {
let token = guard.token();
let _guard = guard;
let outcome = loop {
time::sleep(POLL_INTERVAL).await;
if !SUPERVISOR.monitor_token_valid(&id, token) {
break PollOutcome::TakenOver;
}
let Some(current) = SUPERVISOR.get_daemon(&id).await else {
break PollOutcome::TakenOver;
};
if current.pid == Some(pid) {
let was_stopping = current.status.is_stopping();
PROCS.refresh_pids(&[pid]);
if !PROCS.is_running(pid) {
break PollOutcome::ProcessDied { was_stopping };
}
match PROCS.start_time(pid) {
Some(current_start_time) if current_start_time != expected_start_time => {
break PollOutcome::ProcessDied { was_stopping };
}
None => {
debug!(
"could not read start time for adopted daemon {id} (pid {pid}); continuing on liveness only"
);
}
_ => {}
}
} else if current.pid.is_none() && current.status.is_stopping() {
} else if current.pid.is_none() && current.status.is_stopped() {
break PollOutcome::StoppedExternally;
} else {
break PollOutcome::TakenOver;
}
};
if matches!(outcome, PollOutcome::TakenOver) {
debug!("adopted daemon {id} was taken over or removed; poll monitor exiting");
return;
}
SUPERVISOR
.active_monitors
.fetch_add(1, atomic::Ordering::Release);
struct MonitorGuard;
impl Drop for MonitorGuard {
fn drop(&mut self) {
SUPERVISOR
.active_monitors
.fetch_sub(1, atomic::Ordering::Release);
SUPERVISOR.monitor_done.notify_waiters();
}
}
let _monitor_guard = MonitorGuard;
if !SUPERVISOR.monitor_token_valid(&id, token) {
debug!("adopted daemon {id} was superseded; skipping exit handling");
return;
}
let current = SUPERVISOR.get_daemon(&id).await;
let owns_pid = current.as_ref().is_some_and(|d| d.pid == Some(pid));
let finalized_ours = current.as_ref().is_some_and(|d| {
d.pid.is_none() && (d.status.is_stopped() || d.status.is_stopping())
});
if !owns_pid && !finalized_ours {
debug!("adopted daemon {id} has a successor; skipping exit handling");
return;
}
let was_stopping = matches!(outcome, PollOutcome::ProcessDied { was_stopping: true });
let intentional = was_stopping
|| finalized_ours
|| matches!(outcome, PollOutcome::StoppedExternally)
|| current
.as_ref()
.is_some_and(|d| d.status.is_stopping() || d.status.is_stopped());
let (exit_code, exit_reason) = if intentional {
(-1, "stop")
} else {
(-1, "fail")
};
info!("adopted daemon {id} (pid {pid}) exited ({exit_reason}, exit status unknown)");
if owns_pid {
let new_status = match exit_reason {
"stop" => DaemonStatus::Stopped,
_ => DaemonStatus::Errored(exit_code),
};
let last_exit_success = exit_reason == "stop";
if !SUPERVISOR
.finalize_monitored_exit(&id, pid, token, new_status, Some(last_exit_success))
.await
{
return;
}
}
let hook_extra_env = vec![
("PITCHFORK_EXIT_CODE".to_string(), exit_code.to_string()),
("PITCHFORK_EXIT_REASON".to_string(), exit_reason.to_string()),
];
let hooks_to_fire: Vec<HookType> = match exit_reason {
"stop" => vec![HookType::OnStop, HookType::OnExit],
_ if hook_retry_count >= hook_retry.count() => {
vec![HookType::OnFail, HookType::OnExit]
}
_ => vec![],
};
for hook_type in hooks_to_fire {
fire_hook(
hook_type,
id.clone(),
daemon_dir.clone(),
hook_retry_count,
hook_env.clone(),
hook_extra_env.clone(),
)
.await;
}
});
}
}