use super::{SUPERVISOR, Supervisor};
use crate::Result;
use crate::daemon_id::DaemonId;
use crate::ipc::IpcResponse;
use crate::pitchfork_toml::PitchforkToml;
use crate::settings::settings;
use log::LevelFilter::Info;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use tokio::time;
pub(super) fn is_within(base: &Path, path: &Path) -> bool {
if path.starts_with(base) {
return true;
}
canonicalize_pair(base, path).is_some_and(|(b, p)| p.starts_with(&b))
}
fn dirs_overlap(a: &Path, b: &Path) -> bool {
is_within(a, b) || is_within(b, a)
}
fn canonicalize_pair(a: &Path, b: &Path) -> Option<(PathBuf, PathBuf)> {
Some((a.canonicalize().ok()?, b.canonicalize().ok()?))
}
impl Supervisor {
pub(crate) async fn leave_dir(&self, dir: &Path) -> Result<()> {
debug!("left dir {}", dir.display());
let active_dirs = self.get_active_directories().await;
debug!("active directories after leaving {dir:?}: {active_dirs:?}");
let autostop_delay = settings().general_autostop_delay();
for daemon in self.active_daemons().await {
if !daemon.autostop {
continue;
}
if let Some(daemon_dir) = daemon.dir.as_ref() {
let starts = is_within(dir, daemon_dir);
let still_active = active_dirs.iter().any(|d| is_within(daemon_dir, d));
debug!(
"leave_dir daemon={} daemon_dir={daemon_dir:?} starts_with_left={starts} still_active={still_active}",
daemon.id
);
if starts && !still_active {
if autostop_delay.is_zero() {
info!("autostopping {daemon}");
self.spawn_autostop(daemon.id.clone(), daemon.to_string())
.await;
} else {
let stop_at = time::Instant::now() + autostop_delay;
let mut pending = self.pending_autostops.lock().await;
if !pending.contains_key(&daemon.id) {
info!(
"scheduling autostop for {} in {:?}",
daemon.id, autostop_delay
);
pending.insert(daemon.id.clone(), stop_at);
}
}
}
}
}
Ok(())
}
pub(crate) async fn cancel_pending_autostops_for_dir(&self, dir: &Path) {
let mut pending = self.pending_autostops.lock().await;
let daemons_to_cancel: Vec<DaemonId> = {
let state_file = self.state_file.lock().await;
state_file
.daemons
.iter()
.filter(|(_id, d)| {
d.dir.as_ref().is_some_and(|daemon_dir| {
dirs_overlap(dir, daemon_dir)
})
})
.map(|(id, _)| id.clone())
.collect()
};
for daemon_id in daemons_to_cancel {
if pending.remove(&daemon_id).is_some() {
info!("cancelled pending autostop for {daemon_id}");
}
if let Some(flag) = self.in_flight_autostops.lock().await.remove(&daemon_id) {
flag.store(true, Ordering::SeqCst);
info!("cancelled in-flight autostop for {daemon_id}");
}
}
}
async fn spawn_autostop(&self, daemon_id: DaemonId, label: String) {
let cancelled = Arc::new(AtomicBool::new(false));
{
let mut in_flight = self.in_flight_autostops.lock().await;
if in_flight.contains_key(&daemon_id) {
debug!("autostop already in flight for {daemon_id}");
return;
}
in_flight.insert(daemon_id.clone(), cancelled.clone());
}
tokio::spawn(async move {
SUPERVISOR
.run_autostop(&daemon_id, &label, &cancelled)
.await;
let mut in_flight = SUPERVISOR.in_flight_autostops.lock().await;
if in_flight
.get(&daemon_id)
.is_some_and(|f| Arc::ptr_eq(f, &cancelled))
{
in_flight.remove(&daemon_id);
}
});
}
async fn run_autostop(&self, daemon_id: &DaemonId, label: &str, cancelled: &AtomicBool) {
if cancelled.load(Ordering::SeqCst) {
debug!("autostop of {daemon_id} cancelled before stopping");
return;
}
if let Some(daemon) = self.get_daemon(daemon_id).await
&& daemon.autostop
&& daemon.status.is_running()
{
if let Some(daemon_dir) = daemon.dir.as_ref() {
let active_dirs = self.get_active_directories().await;
if active_dirs.iter().any(|d| is_within(daemon_dir, d)) {
debug!("autostop of {daemon_id} skipped: directory active again");
return;
}
}
} else {
debug!("autostop of {daemon_id} skipped: no longer running");
return;
}
if cancelled.load(Ordering::SeqCst) {
debug!("autostop of {daemon_id} cancelled before stopping");
return;
}
match self.stop(daemon_id).await {
Ok(_) => {
self.add_notification(Info, format!("autostopped {label}"))
.await;
}
Err(e) => error!("failed to autostop {daemon_id}: {e}"),
}
}
pub(crate) async fn process_pending_autostops(&self) -> Result<()> {
let now = time::Instant::now();
let to_stop: Vec<DaemonId> = {
let pending = self.pending_autostops.lock().await;
pending
.iter()
.filter(|(_, stop_at)| now >= **stop_at)
.map(|(id, _)| id.clone())
.collect()
};
for daemon_id in to_stop {
{
let mut pending = self.pending_autostops.lock().await;
pending.remove(&daemon_id);
}
if let Some(daemon) = self.get_daemon(&daemon_id).await
&& daemon.autostop
&& daemon.status.is_running()
{
let active_dirs = self.get_active_directories().await;
if let Some(daemon_dir) = daemon.dir.as_ref() {
let still_active = active_dirs.iter().any(|d| is_within(daemon_dir, d));
debug!(
"process_pending_autostops daemon={daemon_id} daemon_dir={daemon_dir:?} active_dirs={active_dirs:?} still_active={still_active}"
);
if still_active {
debug!(
"process_pending_autostops: daemon={daemon_id} still has active directory, skipping"
);
continue;
}
info!("autostopping {daemon_id} (after delay)");
let label = daemon_id.to_string();
self.spawn_autostop(daemon_id, label).await;
}
}
}
Ok(())
}
pub(crate) async fn start_boot_daemons(&self) -> Result<()> {
info!("Scanning for boot_start daemons");
let pt = PitchforkToml::all_merged_all_namespaces()?;
let boot_daemons: Vec<_> = pt
.daemons
.iter()
.filter(|(_id, d)| d.boot_start.unwrap_or(false))
.collect();
if boot_daemons.is_empty() {
info!("No daemons configured with boot_start = true");
return Ok(());
}
info!("Found {} daemon(s) to start at boot", boot_daemons.len());
for (id, daemon) in boot_daemons {
info!("Starting boot daemon: {id}");
let cmd = match shell_words::split(&daemon.run) {
Ok(cmd) => cmd,
Err(e) => {
error!("failed to parse command for boot daemon {id}: {e}");
continue;
}
};
let mut run_opts = daemon.to_run_options(id, cmd);
run_opts.autostop = false; run_opts.wait_ready = false;
match self.run(run_opts).await {
Ok(IpcResponse::DaemonStart { .. }) | Ok(IpcResponse::DaemonReady { .. }) => {
info!("Successfully started boot daemon: {id}");
}
Ok(IpcResponse::DaemonAlreadyRunning) => {
info!("Boot daemon already running: {id}");
}
Ok(other) => {
warn!("Unexpected response when starting boot daemon {id}: {other:?}");
}
Err(e) => {
error!("Failed to start boot daemon {id}: {e}");
}
}
}
Ok(())
}
}