use crate::Workspace;
use anyhow::Result;
use std::path::Path;
const TEMP_CLEANUP_PROMPT_KEY: &str = "sanitation/temp_cleanup.md";
const TEMP_CLEANUP_WORKSPACE_NAME: &str = "tmp";
pub(crate) async fn temp_cleanup_row_exists(conn: &crate::turso::Connection) -> Result<bool> {
Ok(conn
.query_optional(
"SELECT 1 FROM jobs WHERE kind = 'temp_cleanup' AND status != 'done' LIMIT 1",
(),
|_| Ok::<(), anyhow::Error>(()),
)
.await?
.is_some())
}
pub(crate) async fn dispatch_temp_cleanup() -> Result<()> {
let conn = &crate::session::store().conn;
if temp_cleanup_row_exists(conn).await? {
tracing::info!("Temp-dir cleanup already in flight — skipping dispatch");
return Ok(());
}
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,
)
.await
.map_err(|e| {
tracing::error!(job = %job_id, error = %e, "Failed to spawn temp-dir cleanup job");
e
})?;
let ws = ws.clone();
let job_id_log = job_id.clone();
let agent_id = crate::research_cleanup::cleanup_agent_id(&job_id);
let agent_id_log = agent_id.clone();
tokio::spawn(async move {
run_temp_cleanup_and_finish(&job_id, &ws, &prompt).await;
});
tracing::info!(job = %job_id_log, agent = %agent_id_log, "Temp-dir cleanup dispatched");
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 (agent, response) = crate::agent::run_default_agent(
&agent_id,
crate::Role::Sanitation,
ws,
prompt,
None,
None,
None,
)
.await;
let report = response.unwrap_or_else(|| {
format!(
"Temp-dir cleanup FAILED (job {job_id}): {}",
agent
.failure
.clone()
.unwrap_or_else(|| "no failure detail".to_string())
)
});
tracing::info!(
job = %job_id,
agent = %agent_id,
"Temp-dir cleanup finished: {}",
crate::util::scrub_credentials(&report)
);
let _ = crate::jobs::terminalize_job(&crate::session::store().conn, job_id).await;
}
#[cfg(test)]
mod tests {
use super::*;
async fn init_stores() {
crate::util::test::init_management_test_stores().await;
}
#[tokio::test]
#[serial_test::serial(reset_inflight)] async fn dispatch_temp_cleanup_deduped_by_jobs_row() {
init_stores().await;
crate::jobs::spawn_job(
&crate::session::store().conn,
"tmpclean_dedup",
"task",
"tmp",
"",
"",
crate::Role::Sanitation,
&[],
&crate::jobs::SpawnChild::TempCleanup,
)
.await
.unwrap();
assert!(
temp_cleanup_row_exists(&crate::session::store().conn)
.await
.unwrap(),
"the pre-created row is the dedup marker"
);
dispatch_temp_cleanup().await.unwrap();
let rows = crate::session::store()
.conn
.query("SELECT COUNT(*) FROM jobs WHERE kind = 'temp_cleanup'", ())
.await
.unwrap();
assert_eq!(
rows[0].get::<i64>(0).unwrap(),
1,
"single temp_cleanup job row"
);
let sessions = crate::session::store()
.conn
.query(
"SELECT COUNT(*) FROM session_metadata WHERE agent_id = 'cleanup_tmpclean_dedup'",
(),
)
.await
.unwrap();
assert_eq!(
sessions[0].get::<i64>(0).unwrap(),
0,
"deduped dispatch must not spawn the temp cleaner"
);
crate::jobs::terminalize_job(&crate::session::store().conn, "tmpclean_dedup")
.await
.unwrap();
}
}