use super::*;
use tokio::sync::Mutex as TokioMutex;
fn env_lock() -> std::sync::Arc<TokioMutex<()>> {
static LOCK: std::sync::OnceLock<std::sync::Arc<TokioMutex<()>>> = std::sync::OnceLock::new();
LOCK.get_or_init(|| std::sync::Arc::new(TokioMutex::new(())))
.clone()
}
#[tokio::test]
async fn supervisor_starts_empty() {
let sup = Bm25Supervisor::new();
assert_eq!(sup.supervised_count().await, 0);
}
#[tokio::test]
async fn supervisor_default_matches_new() {
let sup: Bm25Supervisor = Default::default();
assert_eq!(sup.supervised_count().await, 0);
}
#[tokio::test]
async fn external_mode_skips_spawn() {
let lock = env_lock();
let _env = lock.lock().await;
let _guard = EnvGuard::set(ENV_EXTERNAL_BM25, "1");
let tmp = tempfile::tempdir().expect("tempdir");
let sup = Bm25Supervisor::new();
let palace = "ext-skip";
let path = sup
.ensure_running(palace, tmp.path())
.await
.expect("external mode must return socket path without spawning");
assert_eq!(path, socket_path_for_palace(palace));
assert_eq!(
sup.supervised_count().await,
0,
"external mode must not register a child"
);
}
#[tokio::test]
async fn already_running_skips_spawn() {
let lock = env_lock();
let _env = lock.lock().await;
let _g = EnvGuard::remove(ENV_EXTERNAL_BM25);
let palace = format!("a{:x}", std::process::id() & 0xffff);
let socket = socket_path_for_palace(&palace);
let _ = std::fs::remove_file(&socket);
let listener =
trusty_common::uds::bind_hardened(&socket).expect("bind dummy listener at canonical path");
let tmp = tempfile::tempdir().expect("tempdir");
let sup = Bm25Supervisor::new();
let path = sup
.ensure_running(&palace, tmp.path())
.await
.expect("ensure_running must adopt existing socket");
assert_eq!(path, socket);
assert_eq!(
sup.supervised_count().await,
0,
"adoption path must not register a child"
);
drop(listener);
let _ = std::fs::remove_file(&socket);
}
#[tokio::test]
async fn shutdown_with_no_children_is_noop() {
let sup = Bm25Supervisor::new();
sup.shutdown().await;
assert_eq!(sup.supervised_count().await, 0);
}
#[test]
fn supervisor_is_send_and_sync() {
fn assert_send_sync<T: Send + Sync>() {}
assert_send_sync::<Bm25Supervisor>();
}
struct EnvGuard {
key: String,
prev: Option<String>,
}
impl EnvGuard {
fn set(key: &str, value: &str) -> Self {
let prev = std::env::var(key).ok();
unsafe { std::env::set_var(key, value) }
Self {
key: key.to_string(),
prev,
}
}
fn remove(key: &str) -> Self {
let prev = std::env::var(key).ok();
unsafe { std::env::remove_var(key) }
Self {
key: key.to_string(),
prev,
}
}
}
impl Drop for EnvGuard {
fn drop(&mut self) {
unsafe {
match &self.prev {
Some(v) => std::env::set_var(&self.key, v),
None => std::env::remove_var(&self.key),
}
}
}
}
#[test]
fn default_cap_is_three() {
assert_eq!(DEFAULT_MAX_LIVE_DAEMONS, 3);
assert_eq!(Bm25Supervisor::new().max_live(), 3);
}
#[tokio::test]
async fn max_live_daemons_honours_env_override() {
let lock = env_lock();
let _env = lock.lock().await;
let _g = EnvGuard::set(ENV_MAX_DAEMONS, "7");
assert_eq!(max_live_from_env(), 7);
let _g = EnvGuard::set(ENV_MAX_DAEMONS, "0");
assert_eq!(max_live_from_env(), DEFAULT_MAX_LIVE_DAEMONS);
let _g = EnvGuard::set(ENV_MAX_DAEMONS, "banana");
assert_eq!(max_live_from_env(), DEFAULT_MAX_LIVE_DAEMONS);
let _g = EnvGuard::remove(ENV_MAX_DAEMONS);
assert_eq!(max_live_from_env(), DEFAULT_MAX_LIVE_DAEMONS);
}
#[tokio::test]
async fn rss_limit_honours_env_override() {
let lock = env_lock();
let _env = lock.lock().await;
let _g = EnvGuard::set(ENV_RSS_LIMIT_MB, "1024");
assert_eq!(rss_limit_from_env(), Some(1024));
let _g = EnvGuard::set(ENV_RSS_LIMIT_MB, "0");
assert_eq!(rss_limit_from_env(), None, "0 must disable, not reap-all");
let _g = EnvGuard::set(ENV_RSS_LIMIT_MB, "lots");
assert_eq!(rss_limit_from_env(), Some(DEFAULT_RSS_LIMIT_MB));
let _g = EnvGuard::remove(ENV_RSS_LIMIT_MB);
assert_eq!(rss_limit_from_env(), Some(DEFAULT_RSS_LIMIT_MB));
}
#[test]
fn cap_is_clamped_to_at_least_one() {
assert_eq!(Bm25Supervisor::with_limits(0, None).max_live(), 1);
}
#[test]
fn explicit_limits_reach_the_shared_supervisor() {
let sup = Bm25Supervisor::with_limits(9, Some(42));
assert_eq!(sup.max_live(), 9);
assert_eq!(sup.rss_limit_mb(), Some(42));
assert_eq!(
Bm25Supervisor::with_limits(9, None).rss_limit_mb(),
None,
"disabled enforcement must survive the wrapper too"
);
}
#[test]
fn sigterm_patience_exceeds_the_daemon_flush_budget() {
let daemon_flush_budget = trusty_bm25_daemon::SHUTDOWN_FLUSH_TIMEOUT;
assert_eq!(
BM25_TIMEOUTS.shutdown_flush, daemon_flush_budget,
"the supervisor must be configured with the daemon's REAL flush budget, \
or the compile-time guard checks the wrong number"
);
assert!(
BM25_TIMEOUTS.sigterm_patience > daemon_flush_budget,
"the supervisor must outwait the daemon's flush: {:?} vs {daemon_flush_budget:?}",
BM25_TIMEOUTS.sigterm_patience
);
assert!(
BM25_TIMEOUTS.sigterm_patience - daemon_flush_budget >= Duration::from_secs(2),
"leave room for signal delivery, socket cleanup and process exit on top \
of the flush itself"
);
}
#[test]
fn spawn_probe_budget_reaches_the_shared_supervisor() {
assert_eq!(BM25_TIMEOUTS.spawn_probe, Duration::from_millis(3000));
}