use std::path::{Path, PathBuf};
use std::sync::OnceLock;
use std::time::Duration;
use crate::Workspace;
use crate::util::UnwrapPoison;
use anyhow::Result;
use chrono::{DateTime, Duration as ChronoDuration, Utc};
use futures_util::FutureExt;
static TEMP_ROOT: OnceLock<PathBuf> = OnceLock::new();
static LEGACY_TEMP_DIR: OnceLock<PathBuf> = OnceLock::new();
#[must_use]
fn temp_root() -> Option<&'static Path> {
TEMP_ROOT.get().map(PathBuf::as_path)
}
#[must_use]
pub(crate) fn legacy_temp_dir() -> Option<&'static Path> {
LEGACY_TEMP_DIR.get().map(PathBuf::as_path)
}
#[cfg(unix)]
pub fn init_temp_root() -> anyhow::Result<()> {
let legacy = std::env::temp_dir();
let uid = unsafe { libc::geteuid() };
let root = PathBuf::from("/tmp/mahbot");
match std::fs::create_dir(&root) {
Ok(()) => {
std::fs::set_permissions(&root, std::os::unix::fs::PermissionsExt::from_mode(0o700))
.map_err(|e| {
anyhow::anyhow!("temp root {}: chmod 0700 failed: {e}", root.display())
})?;
}
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
use std::os::unix::fs::MetadataExt;
let meta = std::fs::symlink_metadata(&root)
.map_err(|e| anyhow::anyhow!("temp root {}: stat failed: {e}", root.display()))?;
if !meta.is_dir() {
anyhow::bail!(
"temp root {} exists and is not a directory — refusing to use it",
root.display()
);
}
if meta.uid() != uid {
anyhow::bail!(
"temp root {} is owned by uid {} (expected {uid}) — refusing a squatted path",
root.display(),
meta.uid()
);
}
let mode = meta.mode() & 0o777;
if mode & 0o077 != 0 {
std::fs::set_permissions(&root, std::os::unix::fs::PermissionsExt::from_mode(0o700))
.map_err(|e| {
anyhow::anyhow!(
"temp root {} has group/other permissions (mode {mode:o}) and re-chmod 0700 failed: {e}",
root.display()
)
})?;
tracing::warn!(
root = %root.display(),
mode = format_args!("{mode:o}"),
"Temp root had loose permissions — re-chmod 0700 (self-heal, path owned by self)"
);
}
}
Err(e) => anyhow::bail!("temp root {}: create failed: {e}", root.display()),
}
let _ = TEMP_ROOT.set(root.clone());
let _ = LEGACY_TEMP_DIR.set(legacy);
unsafe { std::env::set_var("TMPDIR", &root) };
tracing::info!(root = %root.display(), "Pinned daemon temp root");
Ok(())
}
#[cfg(not(unix))]
pub fn init_temp_root() -> anyhow::Result<()> {
Ok(())
}
#[must_use]
pub(crate) fn shell_tmpdir() -> String {
temp_root().map_or_else(|| "/tmp".to_string(), |p| p.to_string_lossy().into_owned())
}
#[must_use]
pub(crate) fn bare_mktemp_landing_root() -> PathBuf {
#[cfg(target_os = "macos")]
{
if let Some(legacy) = legacy_temp_dir() {
return legacy.to_path_buf();
}
}
std::env::temp_dir()
}
const TEMP_CLEANUP_WAKE: Duration = Duration::from_mins(30);
const GIB: u64 = 1 << 30;
const TEMP_CLEANUP_LOW_FREE_FLOOR: u64 = 10 * GIB;
const TEMP_CLEANUP_RECOVERED_FREE_FLOOR: u64 = 15 * GIB;
const TEMP_CLEANUP_LOW_FREE_PCT: u64 = 10;
const TEMP_CLEANUP_RECOVERED_FREE_PCT: u64 = 15;
pub(crate) const TEMP_CLEANUP_LAST_RUN_KV_KEY: &str = "temp_cleanup_last_run_at";
pub(crate) const TEMP_CLEANUP_MODE_KV_KEY: &str = "temp_cleanup_mode";
static LAST_DISPATCH: std::sync::Mutex<Option<DateTime<Utc>>> = std::sync::Mutex::new(None);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CleanupMode {
Daily,
Weekly,
}
impl CleanupMode {
const fn as_str(self) -> &'static str {
match self {
Self::Daily => "daily",
Self::Weekly => "weekly",
}
}
fn parse(raw: &str) -> Option<Self> {
match raw {
"daily" => Some(Self::Daily),
"weekly" => Some(Self::Weekly),
_ => None,
}
}
fn interval(self) -> ChronoDuration {
match self {
Self::Daily => ChronoDuration::days(1),
Self::Weekly => ChronoDuration::days(7),
}
}
}
fn daily_entry_threshold(capacity: u64) -> u64 {
(capacity.saturating_mul(TEMP_CLEANUP_LOW_FREE_PCT) / 100).max(TEMP_CLEANUP_LOW_FREE_FLOOR)
}
fn weekly_return_threshold(capacity: u64) -> u64 {
(capacity.saturating_mul(TEMP_CLEANUP_RECOVERED_FREE_PCT) / 100)
.max(TEMP_CLEANUP_RECOVERED_FREE_FLOOR)
}
fn mode_after_free(current: CleanupMode, free: u64, capacity: u64) -> CleanupMode {
if free < daily_entry_threshold(capacity) {
CleanupMode::Daily
} else if free > weekly_return_threshold(capacity) {
CleanupMode::Weekly
} else {
current
}
}
async fn stored_cleanup_mode(store: Option<&crate::config_db::ConfigStore>) -> CleanupMode {
let Some(store) = store else {
return CleanupMode::Weekly;
};
match store.get_kv(TEMP_CLEANUP_MODE_KV_KEY).await {
Ok(Some(raw)) => CleanupMode::parse(&raw).unwrap_or_else(|| {
tracing::warn!(value = %raw, "Unrecognized temp cleaner mode — using weekly");
CleanupMode::Weekly
}),
Ok(None) => CleanupMode::Weekly,
Err(e) => {
tracing::warn!(error = %e, "Failed to read temp cleaner mode — using weekly");
CleanupMode::Weekly
}
}
}
fn parse_last_run(raw: &str) -> Option<DateTime<Utc>> {
match crate::db::parse_utc_timestamp(raw) {
Ok(last) => Some(last),
Err(e) => {
tracing::warn!(value = %raw, error = %e, "Unparseable temp cleaner last run — treating as due");
None
}
}
}
fn temp_cleanup_due(
last_run_at: Option<DateTime<Utc>>,
mode: CleanupMode,
now: DateTime<Utc>,
) -> bool {
match last_run_at {
None => true,
Some(last) => now.signed_duration_since(last) >= mode.interval(),
}
}
pub async fn run_temp_cleanup_loop() {
loop {
if !crate::shutdown::sleep_or_shutdown_or_drain(TEMP_CLEANUP_WAKE).await {
break;
}
temp_cleanup_tick().await;
}
}
async fn temp_cleanup_tick() {
let store = crate::config_db::CONFIG_STORE.get();
let current = stored_cleanup_mode(store).await;
let mode = match crate::util::disk::free_and_capacity(&std::env::temp_dir()) {
Some((free, capacity)) => mode_after_free(current, free, capacity),
None => current,
};
if mode != current
&& let Some(store) = store
{
if let Err(e) = store.set_kv(TEMP_CLEANUP_MODE_KV_KEY, mode.as_str()).await {
tracing::warn!(error = %e, "Failed to persist temp cleaner mode");
}
tracing::info!(mode = mode.as_str(), "Temp cleaner cadence mode changed");
}
let persisted = match store {
Some(store) => match store.get_kv(TEMP_CLEANUP_LAST_RUN_KV_KEY).await {
Ok(v) => v,
Err(e) => {
tracing::warn!(error = %e, "Failed to read temp cleaner last run — treating as due");
None
}
},
None => None,
};
let last_run = [
persisted.as_deref().and_then(parse_last_run),
*LAST_DISPATCH.lock().unwrap_poison(),
]
.into_iter()
.flatten()
.max();
if !temp_cleanup_due(last_run, mode, Utc::now()) {
return;
}
if let Err(e) = dispatch_temp_cleanup().await {
tracing::warn!(error = %e, "Temp cleaner dispatch failed");
}
}
const TEMP_CLEANUP_PROMPT_KEY: &str = "sanitation/temp_cleanup.md";
const TEMP_CLEANUP_WORKSPACE_NAME: &str = "tmp";
async fn dispatch_temp_cleanup() -> Result<()> {
let conn = &crate::session::store().conn;
let job_id = crate::generate_id();
let ws = Workspace::ephemeral_run(TEMP_CLEANUP_WORKSPACE_NAME, Path::new("/tmp"));
let prompt = crate::prompt::load_prompt(TEMP_CLEANUP_PROMPT_KEY);
crate::jobs::spawn_job(
conn,
&job_id,
&prompt,
&ws.name,
"",
"",
crate::Role::Sanitation,
&[crate::jobs::NewAgent {
agent_id: crate::research_cleanup::cleanup_agent_id(&job_id),
kind: crate::jobs::AgentKind::Sanitation,
idx: None,
task: prompt.clone(),
}],
&crate::jobs::SpawnChild::TempCleanup,
None,
)
.await
.map_err(|e| {
tracing::error!(job = %job_id, error = %e, "Failed to spawn temp cleaner job");
e
})?;
let now = Utc::now();
*LAST_DISPATCH.lock().unwrap_poison() = Some(now);
if let Some(store) = crate::config_db::CONFIG_STORE.get()
&& let Err(e) = store
.set_kv(TEMP_CLEANUP_LAST_RUN_KV_KEY, &now.to_rfc3339())
.await
{
tracing::warn!(error = %e, "Failed to record temp cleaner run start");
}
tracing::info!(job = %job_id, "Temp cleaner dispatched");
run_temp_cleanup_and_finish(&job_id, &ws, &prompt).await;
Ok(())
}
async fn run_temp_cleanup_and_finish(job_id: &str, ws: &Workspace, prompt: &str) {
let agent_id = crate::research_cleanup::cleanup_agent_id(job_id);
let run = std::panic::AssertUnwindSafe(crate::agent::run_default_agent(
&agent_id,
crate::Role::Sanitation,
ws,
prompt,
false,
None,
None,
None,
))
.catch_unwind()
.await;
match run {
Ok((agent, response)) => {
let report = response.unwrap_or_else(|| {
format!(
"Temp cleaner FAILED (job {job_id}): {}",
agent
.failure
.clone()
.unwrap_or_else(|| "no failure detail".to_string())
)
});
tracing::info!(
job = %job_id,
agent = %agent_id,
"Temp cleaner finished: {}",
crate::util::scrub_credentials(&report)
);
}
Err(payload) => tracing::error!(
job = %job_id,
agent = %agent_id,
"Temp cleaner panicked: {}",
crate::util::panic_message(&*payload)
),
}
let _ = crate::jobs::terminalize_job(&crate::session::store().conn, job_id).await;
}
#[cfg(test)]
mod tests {
use super::*;
#[cfg(unix)]
#[test]
fn legacy_capture_and_pin_are_consistent() {
assert!(temp_root().is_none());
assert!(legacy_temp_dir().is_none());
assert_eq!(shell_tmpdir(), "/tmp");
}
#[test]
fn mode_after_free_applies_hysteresis() {
let small = 50 * GIB;
let big = 460 * GIB;
assert_eq!(daily_entry_threshold(small), 10 * GIB);
assert_eq!(weekly_return_threshold(small), 15 * GIB);
assert_eq!(daily_entry_threshold(big), 46 * GIB);
assert_eq!(weekly_return_threshold(big), 69 * GIB);
assert_eq!(
mode_after_free(CleanupMode::Weekly, 5 * GIB, small),
CleanupMode::Daily
);
assert_eq!(
mode_after_free(CleanupMode::Daily, 20 * GIB, small),
CleanupMode::Weekly
);
assert_eq!(
mode_after_free(CleanupMode::Daily, 12 * GIB, small),
CleanupMode::Daily
);
assert_eq!(
mode_after_free(CleanupMode::Weekly, 12 * GIB, small),
CleanupMode::Weekly
);
}
#[test]
fn due_uses_wall_clock_and_mode_interval() {
let now = Utc::now();
assert!(temp_cleanup_due(None, CleanupMode::Weekly, now));
let two_days_ago = now - ChronoDuration::days(2);
assert!(temp_cleanup_due(
Some(two_days_ago),
CleanupMode::Daily,
now
));
assert!(!temp_cleanup_due(
Some(two_days_ago),
CleanupMode::Weekly,
now
));
let eight_days_ago = now - ChronoDuration::days(8);
assert!(temp_cleanup_due(
Some(eight_days_ago),
CleanupMode::Weekly,
now
));
let recent = now - ChronoDuration::hours(1);
assert!(!temp_cleanup_due(Some(recent), CleanupMode::Daily, now));
}
}