use super::autostop::is_within;
use super::{SUPERVISOR, Supervisor};
use crate::daemon::Daemon;
use crate::daemon_id::DaemonId;
use crate::daemon_status::DaemonStatus;
use crate::ipc::IpcResponse;
use crate::proxy::activity::ACTIVITY;
use log::LevelFilter::Info;
use std::collections::{BTreeMap, HashMap, HashSet};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
static SHELL_ADMISSION: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
pub(crate) async fn admit_shell() -> tokio::sync::MutexGuard<'static, ()> {
SHELL_ADMISSION.lock().await
}
static SWEEPING: AtomicBool = AtomicBool::new(false);
struct Sweep;
impl Sweep {
fn begin() -> Option<Self> {
(!SWEEPING.swap(true, Ordering::SeqCst)).then_some(Sweep)
}
}
impl Drop for Sweep {
fn drop(&mut self) {
SWEEPING.store(false, Ordering::SeqCst);
}
}
fn grace(daemon: &Daemon) -> Option<Duration> {
daemon.proxy_idle_timeout_ms.map(Duration::from_millis)
}
fn is_live(daemon: &Daemon) -> bool {
let status = &daemon.status;
status.is_running()
|| status.is_waiting()
|| status.is_stopping()
|| (status.is_errored() && daemon.retry_count < daemon.retry.count())
}
fn shell_inside(daemon: &Daemon, active_dirs: &[PathBuf]) -> bool {
daemon
.dir
.as_deref()
.is_some_and(|dir| active_dirs.iter().any(|d| is_within(dir, d)))
}
fn live_dependents<'a>(
id: &'a DaemonId,
daemons: &'a BTreeMap<DaemonId, Daemon>,
excluding: &'a HashSet<DaemonId>,
) -> impl Iterator<Item = &'a DaemonId> + 'a {
daemons
.values()
.filter(move |d| is_live(d) && d.depends.contains(id) && !excluding.contains(&d.id))
.map(|d| &d.id)
}
pub(crate) fn plan_idle_stops(
daemons: &BTreeMap<DaemonId, Daemon>,
active_dirs: &[PathBuf],
is_idle: impl Fn(&DaemonId, Duration) -> bool,
) -> Vec<Vec<DaemonId>> {
let mut chosen: HashSet<DaemonId> = daemons
.values()
.filter(|d| d.status.is_running() && !shell_inside(d, active_dirs))
.filter(|d| grace(d).is_some_and(|g| is_idle(&d.id, g)))
.map(|d| d.id.clone())
.collect();
loop {
let needed: Vec<DaemonId> = chosen
.iter()
.filter(|id| live_dependents(id, daemons, &chosen).next().is_some())
.cloned()
.collect();
if needed.is_empty() {
break;
}
for id in needed {
chosen.remove(&id);
}
}
let mut levels = Vec::new();
while !chosen.is_empty() {
let mut level: Vec<DaemonId> = chosen
.iter()
.filter(|id| {
!chosen
.iter()
.any(|other| daemons.get(other).is_some_and(|d| d.depends.contains(id)))
})
.cloned()
.collect();
if level.is_empty() {
break;
}
level.sort();
for id in &level {
chosen.remove(id);
}
levels.push(level);
}
levels
}
impl Supervisor {
pub(crate) async fn claim_daemons(&self, ids: &[DaemonId]) {
let claimed: Vec<DaemonId> = {
let mut state_file = self.state_file.lock().await;
ids.iter()
.filter(|id| state_file.clear_proxy_idle_timeout(id))
.cloned()
.collect()
};
for id in &claimed {
info!("{id} was started explicitly; it will no longer be stopped when idle");
}
for id in ids {
if ACTIVITY.is_idle_stopping(id) {
drop(self.stop_lock(id).await.lock().await);
}
}
}
pub(crate) async fn check_idle_daemons(&self) {
let daemons = {
let state_file = self.state_file.lock().await;
if !state_file
.daemons
.values()
.any(|d| d.proxy_idle_timeout_ms.is_some() && d.status.is_running())
{
return;
}
state_file.daemons.clone()
};
let Some(sweep) = Sweep::begin() else {
return;
};
let active_dirs = self.get_active_directories().await;
let plan = plan_idle_stops(&daemons, &active_dirs, |id, grace| {
let activity = ACTIVITY.snapshot(id);
activity.in_flight == 0 && !activity.idle_stopping && activity.idle_for >= grace
});
if plan.is_empty() {
return;
}
debug!("idle shutdown plan: {plan:?}");
let graces: HashMap<DaemonId, Duration> = daemons
.values()
.filter_map(|d| grace(d).map(|g| (d.id.clone(), g)))
.collect();
tokio::spawn(async move {
let _sweep = sweep;
for level in plan {
SUPERVISOR.idle_stop_level(level, &graces).await;
}
});
}
async fn idle_stop_level(&self, level: Vec<DaemonId>, graces: &HashMap<DaemonId, Duration>) {
let mut group: HashSet<DaemonId> = HashSet::new();
for id in level {
match graces.get(&id) {
Some(&grace) if ACTIVITY.claim_idle_stop(&id, grace) => {
group.insert(id);
}
_ => debug!("idle stop of {id} called off: it was active again"),
}
}
loop {
let mut blocked = Vec::new();
for id in &group {
if let Some(reason) = self.idle_stop_blocker(id, &group, false).await {
debug!("idle stop of {id} called off: {reason}");
blocked.push(id.clone());
}
}
if blocked.is_empty() {
break;
}
for id in blocked {
group.remove(&id);
ACTIVITY.release_idle_stop(&id);
}
}
let mut ordered: Vec<DaemonId> = group.iter().cloned().collect();
ordered.sort();
for id in ordered {
let lock = self.stop_lock(&id).await;
let stopped = {
let _guard = lock.lock().await;
let blocker = {
let _admission = SHELL_ADMISSION.lock().await;
self.idle_stop_blocker(&id, &group, true).await
};
match blocker {
Some(reason) => {
debug!("idle stop of {id} called off: {reason}");
false
}
None => {
let grace = graces.get(&id).copied().unwrap_or_default();
info!("stopping {id}: no proxy activity for {grace:?}");
match self.stop_locked(&id).await {
Ok(IpcResponse::Ok | IpcResponse::DaemonWasNotRunning) => true,
Ok(IpcResponse::DaemonNotRunning) => {
self.settle_stopping(&id, DaemonStatus::Stopped).await;
true
}
Ok(rsp) => {
error!("failed to stop idle daemon {id}: {rsp:?}");
self.settle_stopping(&id, DaemonStatus::Running).await;
false
}
Err(e) => {
error!("failed to stop idle daemon {id}: {e}");
self.settle_stopping(&id, DaemonStatus::Running).await;
false
}
}
}
}
};
ACTIVITY.release_idle_stop(&id);
if stopped {
self.add_notification(Info, format!("stopped idle {id}"))
.await;
} else {
group.remove(&id);
}
}
}
async fn settle_stopping(&self, id: &DaemonId, status: DaemonStatus) {
let mut state_file = self.state_file.lock().await;
if state_file
.daemons
.get(id)
.is_some_and(|d| d.status.is_stopping())
{
state_file.set_status(id, status);
}
}
async fn idle_stop_blocker(
&self,
id: &DaemonId,
stopping_with: &HashSet<DaemonId>,
commit: bool,
) -> Option<&'static str> {
let mut state_file = self.state_file.lock().await;
let active_dirs = state_file.active_directories();
let Some(daemon) = state_file.daemons.get(id) else {
return Some("it is no longer known");
};
if !daemon.status.is_running() {
return Some("it is no longer running");
}
if daemon.proxy_idle_timeout_ms.is_none() {
return Some("it was started explicitly");
}
if shell_inside(daemon, &active_dirs) {
return Some("a shell is inside its directory");
}
if live_dependents(id, &state_file.daemons, stopping_with)
.next()
.is_some()
{
return Some("a running daemon depends on it");
}
if commit {
state_file.set_status(id, DaemonStatus::Stopping);
}
None
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::daemon_status::DaemonStatus;
fn id(name: &str) -> DaemonId {
DaemonId::new("proj", name)
}
struct Fixture(BTreeMap<DaemonId, Daemon>);
impl Fixture {
fn new() -> Self {
Self(BTreeMap::new())
}
fn add(mut self, name: &str, idle_ms: Option<u64>, depends: &[&str]) -> Self {
self.0.insert(
id(name),
Daemon {
id: id(name),
status: DaemonStatus::Running,
dir: Some(PathBuf::from("/work/proj")),
depends: depends.iter().map(|d| id(d)).collect(),
proxy_idle_timeout_ms: idle_ms,
..Default::default()
},
);
self
}
fn status(mut self, name: &str, status: DaemonStatus) -> Self {
self.0.get_mut(&id(name)).unwrap().status = status;
self
}
fn plan(&self, active_dirs: &[&str], idle: &[&str]) -> Vec<Vec<String>> {
let dirs: Vec<PathBuf> = active_dirs.iter().map(PathBuf::from).collect();
let idle: HashSet<DaemonId> = idle.iter().map(|n| id(n)).collect();
plan_idle_stops(&self.0, &dirs, |d, _| idle.contains(d))
.into_iter()
.map(|l| l.into_iter().map(|d| d.name().to_string()).collect())
.collect()
}
}
const G: Option<u64> = Some(60_000);
#[test]
fn only_proxy_started_daemons_are_eligible() {
let f = Fixture::new().add("web", G, &[]).add("manual", None, &[]);
assert_eq!(f.plan(&[], &["web", "manual"]), vec![vec!["web"]]);
}
#[test]
fn active_daemons_are_kept() {
let f = Fixture::new().add("web", G, &[]);
assert!(f.plan(&[], &[]).is_empty());
}
#[test]
fn dependencies_stop_after_their_dependents() {
let f = Fixture::new()
.add("db", G, &[])
.add("cache", G, &[])
.add("api", G, &["db", "cache"])
.add("web", G, &["api"]);
assert_eq!(
f.plan(&[], &["db", "cache", "api", "web"]),
vec![vec!["web"], vec!["api"], vec!["cache", "db"]]
);
}
#[test]
fn a_dependency_stays_while_a_busy_dependent_needs_it() {
let f = Fixture::new().add("db", G, &[]).add("api", G, &["db"]);
assert!(f.plan(&[], &["db"]).is_empty());
}
#[test]
fn a_shared_dependency_stays_for_an_explicitly_started_consumer() {
let f =
Fixture::new()
.add("db", G, &[])
.add("api", G, &["db"])
.add("worker", None, &["db"]);
assert_eq!(f.plan(&[], &["db", "api"]), vec![vec!["api"]]);
}
#[test]
fn a_starting_dependent_keeps_its_dependency() {
let f = Fixture::new()
.add("db", G, &[])
.add("worker", None, &["db"])
.status("worker", DaemonStatus::Waiting);
assert!(f.plan(&[], &["db"]).is_empty());
}
#[test]
fn a_dependent_on_its_way_back_keeps_its_dependency() {
let f = Fixture::new()
.add("db", G, &[])
.add("api", None, &["db"])
.status("api", DaemonStatus::Stopping);
assert!(f.plan(&[], &["db"]).is_empty());
let mut f = Fixture::new()
.add("db", G, &[])
.add("api", None, &["db"])
.status("api", DaemonStatus::Errored(1));
f.0.get_mut(&id("api")).unwrap().retry = crate::config_types::Retry(3);
assert!(f.plan(&[], &["db"]).is_empty());
f.0.get_mut(&id("api")).unwrap().retry_count = 3;
assert_eq!(f.plan(&[], &["db"]), vec![vec!["db"]]);
}
#[test]
fn a_stopped_dependent_does_not_keep_its_dependency() {
let f = Fixture::new()
.add("db", G, &[])
.add("worker", None, &["db"])
.status("worker", DaemonStatus::Stopped);
assert_eq!(f.plan(&[], &["db"]), vec![vec!["db"]]);
}
#[test]
fn a_proxied_daemon_depending_on_a_busy_one_keeps_it() {
let f = Fixture::new().add("api", G, &[]).add("admin", G, &["api"]);
assert!(f.plan(&[], &["api"]).is_empty());
}
#[test]
fn a_shell_inside_the_directory_keeps_the_stack() {
let f = Fixture::new().add("db", G, &[]).add("api", G, &["db"]);
assert!(f.plan(&["/work/proj/src"], &["db", "api"]).is_empty());
assert_eq!(
f.plan(&["/work/other"], &["db", "api"]),
vec![vec!["api"], vec!["db"]]
);
}
#[test]
fn a_dependency_cycle_is_left_running() {
let f = Fixture::new()
.add("db", G, &[])
.add("a", G, &["b", "db"])
.add("b", G, &["a"])
.add("web", G, &["a"]);
assert_eq!(f.plan(&[], &["db", "a", "b", "web"]), vec![vec!["web"]]);
}
}