use std::collections::HashMap;
use std::sync::{LazyLock, Mutex};
use tokio_util::sync::CancellationToken;
use crate::Workspace;
use crate::agent::registry::ParentKey;
use crate::research_cleanup::COMMAND_DUMP_FILE;
use crate::util::UnwrapPoison;
static RESEARCH_CANCELS: LazyLock<ResearchCancelRegistry> =
LazyLock::new(ResearchCancelRegistry::default);
#[derive(Default)]
struct ResearchCancelRegistry {
inner: Mutex<HashMap<String, CancellationToken>>,
}
impl ResearchCancelRegistry {
fn register(&'static self, job_id: &str) -> ResearchCancelGuard {
let token = CancellationToken::new();
self.inner
.lock()
.unwrap_poison()
.insert(job_id.to_string(), token.clone());
ResearchCancelGuard {
job_id: job_id.to_string(),
registry: self,
}
}
fn cancel(&self, job_id: &str) {
let map = self.inner.lock().unwrap_poison();
if let Some(token) = map.get(job_id) {
token.cancel();
}
}
fn is_cancelled(&self, job_id: &str) -> bool {
self.inner
.lock()
.unwrap_poison()
.get(job_id)
.is_some_and(CancellationToken::is_cancelled)
}
fn is_registered(&self, job_id: &str) -> bool {
self.inner.lock().unwrap_poison().contains_key(job_id)
}
fn unregister(&self, job_id: &str) {
self.inner.lock().unwrap_poison().remove(job_id);
}
}
pub(crate) struct ResearchCancelGuard {
job_id: String,
registry: &'static ResearchCancelRegistry,
}
impl Drop for ResearchCancelGuard {
fn drop(&mut self) {
self.registry.unregister(&self.job_id);
}
}
pub(crate) fn register(job_id: &str) -> ResearchCancelGuard {
RESEARCH_CANCELS.register(job_id)
}
pub(crate) fn cancel(job_id: &str) {
RESEARCH_CANCELS.cancel(job_id);
}
pub(crate) fn is_cancelled(job_id: &str) -> bool {
RESEARCH_CANCELS.is_cancelled(job_id)
}
pub(crate) async fn cancel_research_run(job_id: &str) {
cancel(job_id);
crate::agent::registry::AGENT_REGISTRY
.cancel_by_parent_key(&ParentKey::Research(job_id.to_string()));
crate::agent::registry::NON_AGENT_CALLS
.remove_by_parent_key(&ParentKey::Research(job_id.to_string()));
if !RESEARCH_CANCELS.is_registered(job_id)
&& let Err(e) = sweep_cancelled_run(job_id).await
{
tracing::warn!(job = %job_id, error = %e, "cancel sweep failed — durable rows left for boot resume");
}
}
pub(crate) async fn sweep_cancelled_run(job_id: &str) -> Result<(), String> {
let conn = &crate::session::store().conn;
let tx = conn.begin_tx().await.map_err(|e| format!("{e:#}"))?;
tx.execute(
"DELETE FROM pending_jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.map_err(|e| format!("{e:#}"))?;
tx.execute("DELETE FROM jobs WHERE id = ?1", crate::db::params![job_id])
.await
.map_err(|e| format!("{e:#}"))?;
tx.commit().await.map_err(|e| format!("{e:#}"))?;
if crate::research_cleanup::command_dump_exists(job_id).await {
tracing::warn!(
job = %job_id,
"run folder has a command dump — cleanup intent present, folder NOT released (left for the cleanup tail / OS sweep)"
);
} else {
crate::research_cleanup::release_run_folder(job_id).await;
}
delete_results_archive(job_id).await;
Ok(())
}
pub(crate) async fn hand_off_cancelled_run(
job_id: &str,
ws: &Workspace,
question: &str,
) -> Result<(), String> {
if let Some(prompt) = hand_off_rows(job_id, ws).await? {
crate::research_cleanup::spawn_cleanup_agent(job_id, question, ws, &prompt);
}
Ok(())
}
async fn hand_off_rows(job_id: &str, ws: &Workspace) -> Result<Option<String>, String> {
let conn = &crate::session::store().conn;
if crate::research_cleanup::research_cleanup_row_exists(conn, job_id)
.await
.map_err(|e| format!("{e:#}"))?
{
if let Err(e) = conn
.execute(
"DELETE FROM pending_jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
{
tracing::warn!(job = %job_id, error = %e, "pending_jobs removal failed during cancel handoff");
}
delete_results_archive(job_id).await;
return Ok(None);
}
if !crate::research_cleanup::command_dump_exists(job_id).await {
sweep_cancelled_run(job_id).await?;
return Ok(None);
}
let run_root = tokio::fs::canonicalize(crate::research_cleanup::run_root_path(job_id))
.await
.map_err(|e| format!("{e:#}"))?;
let dump_path = run_root.join(COMMAND_DUMP_FILE);
let prompt = crate::research_cleanup::build_cleanup_prompt(job_id, &run_root, &dump_path, ws);
crate::jobs::transition_research_to_cleanup(conn, job_id, &prompt, &ws.name)
.await
.map_err(|e| format!("{e:#}"))?;
delete_results_archive(job_id).await;
Ok(Some(prompt))
}
async fn delete_results_archive(job_id: &str) {
let Some(path) = crate::research_cleanup::results_archive_path(job_id) else {
return;
};
match tokio::fs::remove_file(&path).await {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => {
tracing::warn!(job = %job_id, error = %e, "results.md archive deletion failed — left on disk");
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::util::test::JobRowBuilder;
#[tokio::test]
async fn cancel_fires_signal_and_guard_removes_entry_on_drop() {
cancel("nope");
assert!(!is_cancelled("nope"));
let guard = register("run_signal_1");
assert!(!is_cancelled("run_signal_1"));
cancel("run_signal_1");
assert!(is_cancelled("run_signal_1"), "fired signal stays visible");
cancel("run_signal_1");
assert!(is_cancelled("run_signal_1"));
drop(guard);
assert!(
!is_cancelled("run_signal_1"),
"guard drop removes the entry"
);
}
#[tokio::test]
async fn sweep_removes_rows_folder_and_archive_for_no_dump_run() {
crate::util::test::init_management_test_stores().await;
let job_id = "research_cancel_midrun_1";
let conn = &crate::session::store().conn;
let now = crate::db::now();
JobRowBuilder::new(conn, job_id, "research", "assistant", "ws")
.task("q")
.user_name("u")
.channel("telegram")
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO research_jobs (id, state) VALUES (?1, '{}')",
crate::db::params![job_id],
)
.await
.unwrap();
conn.execute(
"INSERT INTO pending_jobs (id, envelope, target_agent_id, created_at) \
VALUES (?1, ?2, 'manager_ws', ?3)",
crate::db::params![
job_id,
r#"{"content":"report","workspace_name":"ws","user_name":"u","channel":"telegram","kind":"ResearchResult","role":"assistant","reply_target":null,"pending_job_id":null}"#,
now
],
)
.await
.unwrap();
let run_root = crate::research_cleanup::ensure_run_root(job_id).await;
tokio::fs::write(run_root.join("scratch.txt"), "scratch")
.await
.unwrap();
let archive = crate::research_cleanup::results_archive_path(job_id).unwrap();
tokio::fs::create_dir_all(archive.parent().unwrap())
.await
.unwrap();
tokio::fs::write(&archive, "archive").await.unwrap();
sweep_cancelled_run(job_id).await.unwrap();
let jobs = conn
.query(
"SELECT id FROM jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(jobs.is_empty(), "jobs row removed");
let pending = conn
.query(
"SELECT id FROM pending_jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(pending.is_empty(), "pending row removed — no boot replay");
assert!(!run_root.exists(), "run folder removed");
assert!(!archive.exists(), "results.md archive removed");
}
#[tokio::test]
async fn sweep_removes_pending_and_cleanup_rows_for_terminalized_run() {
crate::util::test::init_management_test_stores().await;
let job_id = "research_cancel_terminal_1";
let conn = &crate::session::store().conn;
let now = crate::db::now();
JobRowBuilder::new(conn, job_id, "research_cleanup", "sanitation", "ws")
.task("cleanup prompt")
.user_name("")
.channel("")
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO pending_jobs (id, envelope, target_agent_id, created_at) \
VALUES (?1, ?2, 'manager_ws', ?3)",
crate::db::params![
job_id,
r#"{"content":"report","workspace_name":"ws","user_name":"u","channel":"telegram","kind":"ResearchResult","role":"assistant","reply_target":null,"pending_job_id":null}"#,
now
],
)
.await
.unwrap();
let run_root = crate::research_cleanup::ensure_run_root(job_id).await;
let archive = crate::research_cleanup::results_archive_path(job_id).unwrap();
tokio::fs::create_dir_all(archive.parent().unwrap())
.await
.unwrap();
tokio::fs::write(&archive, "archive").await.unwrap();
sweep_cancelled_run(job_id).await.unwrap();
let jobs = conn
.query(
"SELECT id FROM jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(jobs.is_empty(), "cleanup job row removed");
let pending = conn
.query(
"SELECT id FROM pending_jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(pending.is_empty(), "pending report row removed");
assert!(!run_root.exists());
assert!(!archive.exists());
}
#[tokio::test]
async fn sweep_never_releases_folder_with_command_dump() {
crate::util::test::init_management_test_stores().await;
let job_id = "research_cancel_dump_guard_1";
let conn = &crate::session::store().conn;
let now = crate::db::now();
JobRowBuilder::new(conn, job_id, "research", "assistant", "ws")
.task("q")
.user_name("u")
.channel("telegram")
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO pending_jobs (id, envelope, target_agent_id, created_at) \
VALUES (?1, ?2, 'manager_ws', ?3)",
crate::db::params![
job_id,
r#"{"content":"report","workspace_name":"ws","user_name":"u","channel":"telegram","kind":"ResearchResult","role":"assistant","reply_target":null,"pending_job_id":null}"#,
now
],
)
.await
.unwrap();
let run_root = crate::research_cleanup::ensure_run_root(job_id).await;
tokio::fs::write(run_root.join(COMMAND_DUMP_FILE), "shell cmd")
.await
.unwrap();
let archive = crate::research_cleanup::results_archive_path(job_id).unwrap();
tokio::fs::create_dir_all(archive.parent().unwrap())
.await
.unwrap();
tokio::fs::write(&archive, "archive").await.unwrap();
sweep_cancelled_run(job_id).await.unwrap();
let jobs = conn
.query(
"SELECT id FROM jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(jobs.is_empty(), "jobs row removed");
let pending = conn
.query(
"SELECT id FROM pending_jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(pending.is_empty(), "pending row removed");
assert!(
run_root.exists(),
"run folder held — cleanup intent present"
);
assert!(
run_root.join(COMMAND_DUMP_FILE).exists(),
"command dump survives"
);
assert!(!archive.exists(), "results.md archive removed");
}
#[tokio::test]
#[expect(clippy::too_many_lines)] async fn cancel_handoff_transitions_rows_and_holds_folder() {
crate::util::test::init_management_test_stores().await;
let job_id = "research_handoff_transition_1";
let conn = &crate::session::store().conn;
let now = crate::db::now();
JobRowBuilder::new(conn, job_id, "research", "assistant", "ws")
.task("question?")
.user_name("u")
.channel("telegram")
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO research_jobs (id, state) VALUES (?1, '{}')",
crate::db::params![job_id],
)
.await
.unwrap();
conn.execute(
"INSERT INTO pending_jobs (id, envelope, target_agent_id, created_at) \
VALUES (?1, ?2, 'manager_ws', ?3)",
crate::db::params![
job_id,
r#"{"content":"report","workspace_name":"ws","user_name":"u","channel":"telegram","kind":"ResearchResult","role":"assistant","reply_target":null,"pending_job_id":null}"#,
now
],
)
.await
.unwrap();
let run_root = crate::research_cleanup::ensure_run_root(job_id).await;
tokio::fs::write(run_root.join(COMMAND_DUMP_FILE), "shell cmd")
.await
.unwrap();
let archive = crate::research_cleanup::results_archive_path(job_id).unwrap();
tokio::fs::create_dir_all(archive.parent().unwrap())
.await
.unwrap();
tokio::fs::write(&archive, "archive").await.unwrap();
let ws = crate::workspace::test_ws("/tmp/test_ws_handoff_transition");
let prompt = hand_off_rows(job_id, &ws).await.unwrap();
let prompt = prompt.expect("handoff transitions to cleanup and returns the prompt");
let row = conn
.query(
"SELECT kind, role, task FROM jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert_eq!(row.len(), 1, "single jobs row after transition");
assert_eq!(
row[0].get::<String>(0).unwrap(),
"research_cleanup",
"kind transitioned to research_cleanup"
);
assert_eq!(
row[0].get::<String>(1).unwrap(),
"sanitation",
"role is sanitation"
);
assert_eq!(
row[0].get::<String>(2).unwrap(),
prompt,
"task is the returned cleanup prompt"
);
let roster = conn
.query(
"SELECT agent_id FROM agents WHERE job_id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert_eq!(roster.len(), 1, "single cleanup roster row");
assert_eq!(
roster[0].get::<String>(0).unwrap(),
format!("cleanup_{job_id}"),
"cleanup agent roster row present"
);
let pending = conn
.query(
"SELECT id FROM pending_jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(pending.is_empty(), "pending envelope dropped");
let child = conn
.query(
"SELECT id FROM research_jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(
child.is_empty(),
"research_jobs child row removed by cascade"
);
assert!(run_root.exists(), "run folder held for the cleanup tail");
assert!(
run_root.join(COMMAND_DUMP_FILE).exists(),
"command dump survives — cleanup intent preserved"
);
assert!(!archive.exists(), "results.md archive removed");
}
#[tokio::test]
async fn cancel_handoff_without_dump_sweeps() {
crate::util::test::init_management_test_stores().await;
let job_id = "research_handoff_no_dump_1";
let conn = &crate::session::store().conn;
let now = crate::db::now();
JobRowBuilder::new(conn, job_id, "research", "assistant", "ws")
.task("question?")
.user_name("u")
.channel("telegram")
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO research_jobs (id, state) VALUES (?1, '{}')",
crate::db::params![job_id],
)
.await
.unwrap();
let run_root = crate::research_cleanup::ensure_run_root(job_id).await;
tokio::fs::write(run_root.join("scratch.txt"), "scratch")
.await
.unwrap();
let ws = crate::workspace::test_ws("/tmp/test_ws_handoff_no_dump");
let outcome = hand_off_rows(job_id, &ws).await.unwrap();
assert!(outcome.is_none(), "no dump → no cleanup prompt");
let jobs = conn
.query(
"SELECT id FROM jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(jobs.is_empty(), "research row swept");
assert!(!run_root.exists(), "dump-less folder released");
}
#[tokio::test]
async fn cancel_research_run_sweeps_unregistered_run() {
crate::util::test::init_management_test_stores().await;
let job_id = "research_cancel_unregistered_1";
let conn = &crate::session::store().conn;
let now = crate::db::now();
JobRowBuilder::new(conn, job_id, "research", "assistant", "ws")
.task("q")
.user_name("u")
.channel("telegram")
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO research_jobs (id, state) VALUES (?1, '{}')",
crate::db::params![job_id],
)
.await
.unwrap();
let run_root = crate::research_cleanup::ensure_run_root(job_id).await;
tokio::fs::write(run_root.join("scratch.txt"), "scratch")
.await
.unwrap();
crate::research_cancel::cancel_research_run(job_id).await;
let jobs = conn
.query(
"SELECT id FROM jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(
jobs.is_empty(),
"research row swept via the unregistered path"
);
assert!(!run_root.exists(), "dump-less folder released");
}
#[tokio::test]
async fn sweep_is_idempotent_and_double_release_safe() {
crate::util::test::init_management_test_stores().await;
let job_id = "research_cancel_double_1";
sweep_cancelled_run(job_id).await.unwrap();
sweep_cancelled_run(job_id).await.unwrap();
}
#[tokio::test]
async fn complete_durable_job_rolls_back_when_run_cancelled() {
crate::util::test::init_management_test_stores().await;
let job_id = "research_cancel_complete_1";
let ws = crate::workspace::test_ws("/tmp/test_ws_research_cancel_complete");
crate::jobs::spawn_job(
&crate::session::store().conn,
job_id,
"q",
&ws.name,
"caller-user",
"telegram",
crate::Role::Assistant,
&[],
&crate::jobs::SpawnChild::Research,
None,
)
.await
.unwrap();
let _guard = register(job_id);
cancel(job_id);
let envelope = crate::jobs::complete_durable_job(
job_id,
"report".to_string(),
crate::agent::message_router::MessageKind::ResearchResult,
crate::Role::Assistant,
"caller-user",
"telegram",
&ws.name,
)
.await;
let conn = &crate::session::store().conn;
let pending = conn
.query(
"SELECT id FROM pending_jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(
pending.is_empty(),
"rolled-back completion leaves no pending row"
);
let jobs = conn
.query(
"SELECT id FROM jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert_eq!(
jobs.len(),
1,
"job row survives the rollback — the sweep deletes it"
);
assert!(
envelope.pending_job_id.is_some(),
"the caller gate (is_cancelled) decides routing"
);
}
}