use std::path::{Path, PathBuf};
use crate::adapter::filesystem::AdmissionLock;
use crate::adapter::git::GitCli;
use crate::adapter::sqlite_repos::{
SqliteAgentRepository, SqliteCleanupRepository, SqliteEventRepository, SqliteSessionRepository,
SqliteWorktreeRepository,
};
use crate::adapter::unit_of_work::UnitOfWork;
use crate::domain::cleanup::CleanupState;
use crate::domain::cleanup::{CleanupAssessment, CleanupBlocker};
use crate::domain::session::SessionState;
use crate::error::CarryCtxError;
use crate::repository::TaskRepository;
use crate::repository::{CleanupRepository, EventRepository, NewEvent};
use crate::repository::{SessionRepository, WorktreeRepository};
fn resolve_reference(
cleanup_repo: &dyn CleanupRepository,
task_repo: &dyn TaskRepository,
project_id: &str,
reference: &str,
) -> Result<crate::repository::CleanupRecord, CarryCtxError> {
if let Some(request) = cleanup_repo.find_by_id(project_id, reference)? {
return Ok(request);
}
let task = match task_repo.find_by_display_id(project_id, reference)? {
Some(task) => task,
None => task_repo
.find_by_id(project_id, reference)?
.ok_or_else(|| CarryCtxError::resource_not_found("Cleanup request not found."))?,
};
let requests = cleanup_repo.find_by_task(project_id, &task.id)?;
requests
.iter()
.find(|request| request.state.is_active() || request.state == CleanupState::Failed)
.cloned()
.or_else(|| requests.into_iter().next())
.ok_or_else(|| CarryCtxError::resource_not_found("Cleanup request not found."))
}
pub fn show_request(
cleanup_repo: &dyn CleanupRepository,
task_repo: &dyn TaskRepository,
project_id: &str,
reference: &str,
) -> Result<crate::repository::CleanupRecord, CarryCtxError> {
resolve_reference(cleanup_repo, task_repo, project_id, reference)
}
pub fn preview_requests(
cleanup_repo: &dyn CleanupRepository,
task_repo: &dyn TaskRepository,
project_id: &str,
reference: Option<&str>,
) -> Result<Vec<crate::repository::CleanupRecord>, CarryCtxError> {
match reference {
Some(reference) => Ok(vec![resolve_reference(
cleanup_repo,
task_repo,
project_id,
reference,
)?]),
None => cleanup_repo.find_pending_by_project(project_id),
}
}
pub fn run_requests(
conn: &mut rusqlite::Connection,
project_id: &str,
reference: Option<&str>,
repo_root: &Path,
actor_agent_id: Option<&str>,
admission_lock: &AdmissionLock,
) -> Result<(Vec<crate::repository::CleanupRecord>, Vec<String>), CarryCtxError> {
run_requests_with_policy(
conn,
project_id,
reference,
repo_root,
actor_agent_id,
admission_lock,
&crate::domain::config::WorktreeCleanupConfig::default(),
)
}
pub fn run_requests_with_policy(
conn: &mut rusqlite::Connection,
project_id: &str,
reference: Option<&str>,
repo_root: &Path,
actor_agent_id: Option<&str>,
admission_lock: &AdmissionLock,
cleanup_config: &crate::domain::config::WorktreeCleanupConfig,
) -> Result<(Vec<crate::repository::CleanupRecord>, Vec<String>), CarryCtxError> {
let cleanup_repo = SqliteCleanupRepository::new(conn);
let task_repo = crate::adapter::sqlite_repos::SqliteTaskRepository::new(conn);
let requests = preview_requests(&cleanup_repo, &task_repo, project_id, reference)?;
let mut warnings = Vec::new();
for request in requests {
warnings.extend(try_cleanup_request_with_policy(
conn,
project_id,
&request.id,
repo_root,
actor_agent_id,
admission_lock,
cleanup_config,
)?);
}
Ok((
SqliteCleanupRepository::new(conn).list(project_id, None)?,
warnings,
))
}
#[derive(Debug)]
pub enum ExecuteOutcome {
Removed,
AlreadyRemoved,
Blocked(CleanupAssessment),
}
pub fn try_cleanup(
conn: &mut rusqlite::Connection,
project_id: &str,
task_id: &str,
repo_root: &Path,
actor_agent_id: Option<&str>,
admission_lock: &AdmissionLock,
) -> Result<Vec<String>, CarryCtxError> {
let request = SqliteCleanupRepository::new(conn)
.find_by_task(project_id, task_id)?
.into_iter()
.find(|r| r.state.is_active() || r.state == CleanupState::Failed);
let Some(request) = request else {
return Ok(Vec::new());
};
reconcile_request_record(
conn,
project_id,
request,
repo_root,
actor_agent_id,
admission_lock,
&crate::domain::config::WorktreeCleanupConfig::default(),
)
}
pub fn try_cleanup_request(
conn: &mut rusqlite::Connection,
project_id: &str,
request_id: &str,
repo_root: &Path,
actor_agent_id: Option<&str>,
admission_lock: &AdmissionLock,
) -> Result<Vec<String>, CarryCtxError> {
try_cleanup_request_with_policy(
conn,
project_id,
request_id,
repo_root,
actor_agent_id,
admission_lock,
&crate::domain::config::WorktreeCleanupConfig::default(),
)
}
pub fn try_cleanup_request_with_policy(
conn: &mut rusqlite::Connection,
project_id: &str,
request_id: &str,
repo_root: &Path,
actor_agent_id: Option<&str>,
admission_lock: &AdmissionLock,
cleanup_config: &crate::domain::config::WorktreeCleanupConfig,
) -> Result<Vec<String>, CarryCtxError> {
let request = SqliteCleanupRepository::new(conn)
.find_by_id(project_id, request_id)?
.ok_or_else(|| CarryCtxError::resource_not_found("Cleanup request not found."))?;
if matches!(
request.state,
CleanupState::Completed | CleanupState::Cancelled
) {
return Err(CarryCtxError::state_conflict(format!(
"Cleanup request '{}' is not retryable (state: {}).",
request.id, request.state
)));
}
if !request.state.is_active() && request.state != CleanupState::Failed {
return Ok(Vec::new());
}
reconcile_request_record(
conn,
project_id,
request,
repo_root,
actor_agent_id,
admission_lock,
cleanup_config,
)
}
fn reconcile_request_record(
conn: &mut rusqlite::Connection,
project_id: &str,
request: crate::repository::CleanupRecord,
repo_root: &Path,
actor_agent_id: Option<&str>,
admission_lock: &AdmissionLock,
cleanup_config: &crate::domain::config::WorktreeCleanupConfig,
) -> Result<Vec<String>, CarryCtxError> {
let _admission_lock = admission_lock;
let now = chrono::Utc::now().to_rfc3339();
let running = {
let uow = UnitOfWork::begin(conn)?;
let repo = SqliteCleanupRepository::new(uow.connection());
let Some(record) = repo.claim_for_attempt(
&request.id,
project_id,
request.state,
request.last_attempt_at.as_deref(),
&now,
)?
else {
uow.rollback()?;
return Ok(Vec::new());
};
let events = SqliteEventRepository::new(uow.connection());
events.append(&NewEvent {
id: ulid::Ulid::generate().to_string(),
project_id: project_id.to_string(),
event_type: "worktree.cleanup_started".into(),
actor_agent_id: actor_agent_id.map(str::to_owned),
session_id: None,
task_id: record.task_id.clone(),
payload: serde_json::json!({
"cleanup_id": record.id,
"reason": record.reason,
"status": record.state,
"attempt_count": record.attempt_count,
}),
occurred_at: now.clone(),
})?;
uow.commit()?;
record
};
let outcome: Result<ExecuteOutcome, CarryCtxError> = {
let sessions = SqliteSessionRepository::new(conn);
let worktrees = SqliteWorktreeRepository::new(conn);
assess_worktree_cleanup_with_policy(
&sessions,
&worktrees,
&GitCli::new(),
project_id,
running.worktree_id.as_deref(),
Path::new(&running.worktree_path),
repo_root,
None,
cleanup_config,
)
.and_then(|assessment| {
if !assessment.removable {
Ok(ExecuteOutcome::Blocked(assessment))
} else {
execute_worktree_cleanup(
&worktrees,
&GitCli::new(),
project_id,
running.worktree_id.as_deref(),
Path::new(&running.worktree_path),
repo_root,
false,
)
}
})
};
let (state, blocker, failure_reason, warning) = match outcome {
Ok(ExecuteOutcome::Removed | ExecuteOutcome::AlreadyRemoved) => {
(CleanupState::Completed, None, None, None)
}
Ok(ExecuteOutcome::Blocked(a)) => {
let blocker = a.blockers.first().cloned();
(
CleanupState::Blocked,
blocker,
None,
Some(format!(
"Worktree cleanup deferred: {}",
a.blockers
.iter()
.map(ToString::to_string)
.collect::<Vec<_>>()
.join(", ")
)),
)
}
Err(err) => (
CleanupState::Failed,
None,
Some(err.message.clone()),
Some(format!("Worktree cleanup failed: {}", err.message)),
),
};
let uow = UnitOfWork::begin(conn)?;
let repo = SqliteCleanupRepository::new(uow.connection());
let record = repo.update_state(
&request.id,
project_id,
state,
blocker.clone(),
&chrono::Utc::now().to_rfc3339(),
)?;
if state == CleanupState::Completed
&& let Some(worktree_id) = request.worktree_id.as_deref()
{
SqliteWorktreeRepository::new(uow.connection()).delete(worktree_id, project_id)?;
}
let actor_agent_id = crate::application::task::canonical_actor_id(
project_id,
actor_agent_id,
&SqliteAgentRepository::new(uow.connection()),
)?;
let events = SqliteEventRepository::new(uow.connection());
events.append(&NewEvent {
id: ulid::Ulid::generate().to_string(),
project_id: project_id.to_string(),
event_type: match state {
CleanupState::Completed => "worktree.removed",
CleanupState::Blocked => "worktree.cleanup_blocked",
CleanupState::Failed => "worktree.cleanup_failed",
_ => unreachable!("try_cleanup only persists completed, blocked, or failed"),
}
.into(),
actor_agent_id,
session_id: None,
task_id: request.task_id.clone(),
payload: serde_json::json!({
"cleanup_id": record.id,
"reason": record.reason,
"status": state,
"attempt_count": record.attempt_count,
"blocked_reason": blocker.as_ref().map(ToString::to_string),
"error": failure_reason,
}),
occurred_at: chrono::Utc::now().to_rfc3339(),
})?;
uow.commit()?;
Ok(warning.into_iter().collect())
}
pub fn reconcile_pending_cleanup(
conn: &mut rusqlite::Connection,
project_id: &str,
repo_root: &Path,
actor_agent_id: Option<&str>,
admission_lock: &AdmissionLock,
) -> Result<Vec<String>, CarryCtxError> {
reconcile_pending_cleanup_with_policy(
conn,
project_id,
repo_root,
actor_agent_id,
admission_lock,
&crate::domain::config::WorktreeCleanupConfig::default(),
)
}
pub fn reconcile_pending_cleanup_with_policy(
conn: &mut rusqlite::Connection,
project_id: &str,
repo_root: &Path,
actor_agent_id: Option<&str>,
admission_lock: &AdmissionLock,
cleanup_config: &crate::domain::config::WorktreeCleanupConfig,
) -> Result<Vec<String>, CarryCtxError> {
let requests = SqliteCleanupRepository::new(conn).find_pending_by_project(project_id)?;
let mut warnings = Vec::new();
for request in requests {
warnings.extend(try_cleanup_request_with_policy(
conn,
project_id,
&request.id,
repo_root,
actor_agent_id,
admission_lock,
cleanup_config,
)?);
}
Ok(warnings)
}
pub fn reconcile_cleanup_for_session(
conn: &mut rusqlite::Connection,
project_id: &str,
task_id: Option<&str>,
worktree_id: Option<&str>,
repo_root: &Path,
actor_agent_id: Option<&str>,
admission_lock: &AdmissionLock,
) -> Result<Vec<String>, CarryCtxError> {
reconcile_cleanup_for_session_with_policy(
conn,
project_id,
task_id,
worktree_id,
repo_root,
actor_agent_id,
admission_lock,
&crate::domain::config::WorktreeCleanupConfig::default(),
)
}
pub fn reconcile_cleanup_for_session_with_policy(
conn: &mut rusqlite::Connection,
project_id: &str,
task_id: Option<&str>,
worktree_id: Option<&str>,
repo_root: &Path,
actor_agent_id: Option<&str>,
admission_lock: &AdmissionLock,
cleanup_config: &crate::domain::config::WorktreeCleanupConfig,
) -> Result<Vec<String>, CarryCtxError> {
let requests = SqliteCleanupRepository::new(conn).find_pending_by_project(project_id)?;
let mut warnings = Vec::new();
for request in requests.into_iter().filter(|request| match worktree_id {
Some(id) => request.worktree_id.as_deref() == Some(id),
None => task_id.is_some_and(|id| request.task_id.as_deref() == Some(id)),
}) {
warnings.extend(try_cleanup_request_with_policy(
conn,
project_id,
&request.id,
repo_root,
actor_agent_id,
admission_lock,
cleanup_config,
)?);
}
Ok(warnings)
}
pub fn assess_worktree_cleanup(
session_repo: &dyn SessionRepository,
worktree_repo: &dyn WorktreeRepository,
git_cli: &GitCli,
project_id: &str,
worktree_id: Option<&str>,
worktree_path: &Path,
repo_root: &Path,
current_dir_override: Option<&Path>,
) -> Result<CleanupAssessment, CarryCtxError> {
assess_worktree_cleanup_with_policy(
session_repo,
worktree_repo,
git_cli,
project_id,
worktree_id,
worktree_path,
repo_root,
current_dir_override,
&crate::domain::config::WorktreeCleanupConfig::default(),
)
}
pub fn assess_worktree_cleanup_with_policy(
session_repo: &dyn SessionRepository,
worktree_repo: &dyn WorktreeRepository,
git_cli: &GitCli,
project_id: &str,
worktree_id: Option<&str>,
worktree_path: &Path,
repo_root: &Path,
current_dir_override: Option<&Path>,
cleanup_config: &crate::domain::config::WorktreeCleanupConfig,
) -> Result<CleanupAssessment, CarryCtxError> {
let mut blockers = Vec::new();
if cleanup_config.require_no_active_session {
if let Some(ids) =
collect_active_session_blockers(session_repo, project_id, worktree_id, worktree_path)?
{
blockers.extend(ids);
}
}
if is_current_dir_inside(worktree_path, current_dir_override) {
blockers.push(CleanupBlocker::CurrentWorkingDirectory);
}
if is_missing_git_metadata(
worktree_repo,
git_cli,
project_id,
worktree_id,
worktree_path,
repo_root,
)? {
blockers.push(CleanupBlocker::MissingGitMetadata);
}
if worktree_path.exists()
&& crate::adapter::git::detect_jj_colocation(&git_cli.discover(repo_root)?.git_common_dir)
{
blockers.push(CleanupBlocker::JjColocation);
}
if cleanup_config.require_clean
&& worktree_path.exists()
&& is_dirty_worktree(git_cli, worktree_path)
{
blockers.push(CleanupBlocker::DirtyWorktree);
}
if is_worktree_locked(git_cli, repo_root, worktree_path) {
blockers.push(CleanupBlocker::WorktreeLocked);
}
if blockers.is_empty() {
Ok(CleanupAssessment::removable())
} else {
Ok(CleanupAssessment::blocked(blockers))
}
}
pub fn assess_worktree_cleanup_simple(
session_repo: &dyn SessionRepository,
worktree_repo: &dyn WorktreeRepository,
git_cli: &GitCli,
project_id: &str,
worktree_id: Option<&str>,
worktree_path: &Path,
repo_root: &Path,
) -> Result<CleanupAssessment, CarryCtxError> {
assess_worktree_cleanup(
session_repo,
worktree_repo,
git_cli,
project_id,
worktree_id,
worktree_path,
repo_root,
None,
)
}
pub fn execute_worktree_cleanup(
worktree_repo: &dyn WorktreeRepository,
git_cli: &GitCli,
project_id: &str,
worktree_id: Option<&str>,
worktree_path: &Path,
repo_root: &Path,
force: bool,
) -> Result<ExecuteOutcome, CarryCtxError> {
let row_missing =
is_worktree_row_missing(worktree_repo, project_id, worktree_id, worktree_path)?;
let path_missing = !worktree_path.exists();
if row_missing || path_missing {
return Ok(ExecuteOutcome::AlreadyRemoved);
}
if crate::adapter::git::detect_jj_colocation(&git_cli.discover(repo_root)?.git_common_dir) {
return Ok(ExecuteOutcome::Blocked(CleanupAssessment::blocked(vec![
CleanupBlocker::JjColocation,
])));
}
match git_cli.remove_worktree(repo_root, worktree_path, force) {
Ok(()) => Ok(ExecuteOutcome::Removed),
Err(err) => {
let msg = err.message.to_lowercase();
if msg.contains("is not a working tree") || msg.contains("not a working tree") {
return Ok(ExecuteOutcome::AlreadyRemoved);
}
if msg.contains("locked working tree") || msg.contains("lock reason") {
let assessment = CleanupAssessment::blocked(vec![CleanupBlocker::WorktreeLocked]);
return Ok(ExecuteOutcome::Blocked(assessment));
}
if err.code == "STATE_CONFLICT" {
let assessment = CleanupAssessment::blocked(vec![CleanupBlocker::DirtyWorktree]);
return Ok(ExecuteOutcome::Blocked(assessment));
}
if msg.contains("current working directory") || msg.contains("is the current") {
let assessment =
CleanupAssessment::blocked(vec![CleanupBlocker::CurrentWorkingDirectory]);
return Ok(ExecuteOutcome::Blocked(assessment));
}
Err(err)
}
}
}
fn collect_active_session_blockers(
session_repo: &dyn SessionRepository,
project_id: &str,
worktree_id: Option<&str>,
worktree_path: &Path,
) -> Result<Option<Vec<CleanupBlocker>>, CarryCtxError> {
let sessions = session_repo.list(project_id)?;
let wp_str = worktree_path.to_string_lossy().to_string();
let mut out = Vec::new();
for s in sessions {
if s.state != SessionState::Active {
continue;
}
let matches_id = match (worktree_id, s.worktree_id.as_deref()) {
(Some(wid), Some(sid)) => wid == sid,
_ => false,
};
let matches_cwd = if let Some(cwd) = s.cwd.as_deref() {
cwd == wp_str
|| Path::new(cwd).starts_with(worktree_path)
|| wp_str.starts_with(cwd) && cwd_matches_worktree(cwd, worktree_path)
} else {
false
};
if matches_id || matches_cwd {
out.push(CleanupBlocker::ActiveSession {
session_id: s.id.clone(),
});
}
}
if out.is_empty() {
Ok(None)
} else {
Ok(Some(out))
}
}
fn cwd_matches_worktree(cwd: &str, worktree_path: &Path) -> bool {
let cwd_path = Path::new(cwd);
cwd_path.starts_with(worktree_path) || worktree_path.starts_with(cwd_path)
}
fn is_current_dir_inside(worktree_path: &Path, override_dir: Option<&Path>) -> bool {
let current = if let Some(p) = override_dir {
p.to_path_buf()
} else {
match std::env::current_dir() {
Ok(p) => p,
Err(_) => return true,
}
};
let wt_canon = worktree_path
.canonicalize()
.unwrap_or_else(|_| worktree_path.to_path_buf());
let cur_canon = current.canonicalize().unwrap_or(current);
cur_canon == wt_canon || cur_canon.starts_with(&wt_canon)
}
fn is_missing_git_metadata(
worktree_repo: &dyn WorktreeRepository,
git_cli: &GitCli,
project_id: &str,
worktree_id: Option<&str>,
worktree_path: &Path,
repo_root: &Path,
) -> Result<bool, CarryCtxError> {
if let Some(wid) = worktree_id {
match worktree_repo.find_by_id(project_id, wid) {
Ok(None) | Err(_) => return Ok(true),
Ok(Some(_)) => {}
}
}
if !worktree_path.exists() {
return Ok(false);
}
let is_git_worktree = match git_cli.list_worktrees(repo_root) {
Ok(entries) => entries.iter().any(|e| {
Path::new(&e.path) == worktree_path
|| same_path_canonical(Path::new(&e.path), worktree_path)
}),
Err(_) => return Ok(true),
};
if is_git_worktree {
return Ok(false);
}
if worktree_id.is_some() {
return Ok(true);
}
let git_file = worktree_path.join(".git");
Ok(git_file.exists())
}
fn same_path_canonical(a: &Path, b: &Path) -> bool {
let a_c = a.canonicalize().unwrap_or_else(|_| a.to_path_buf());
let b_c = b.canonicalize().unwrap_or_else(|_| b.to_path_buf());
a_c == b_c
}
fn is_dirty_worktree(git_cli: &GitCli, worktree_path: &Path) -> bool {
match git_cli.get_snapshot(worktree_path) {
Ok(snap) => snap.dirty,
Err(_) => true,
}
}
fn is_worktree_locked(git_cli: &GitCli, repo_root: &Path, worktree_path: &Path) -> bool {
match git_cli.list_worktrees(repo_root) {
Ok(entries) => {
for e in entries {
let matches = Path::new(&e.path) == worktree_path
|| same_path_canonical(Path::new(&e.path), worktree_path);
if matches {
return e.locked.is_some();
}
}
false
}
Err(_) => {
let _ = filesystem_locked_fallback(repo_root, worktree_path);
true
}
}
}
fn filesystem_locked_fallback(_repo_root: &Path, _worktree_path: &Path) -> bool {
false
}
fn is_worktree_row_missing(
worktree_repo: &dyn WorktreeRepository,
project_id: &str,
worktree_id: Option<&str>,
worktree_path: &Path,
) -> Result<bool, CarryCtxError> {
if let Some(wid) = worktree_id {
if worktree_repo.find_by_id(project_id, wid)?.is_some() {
return Ok(false);
}
}
match worktree_repo.find_by_path(project_id, &worktree_path.to_string_lossy()) {
Ok(Some(_)) => Ok(false),
Ok(None) => Ok(true),
Err(error) => Err(error),
}
}
#[allow(dead_code)]
fn _use_pathbuf(p: PathBuf) -> PathBuf {
p
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapter::git::GitCli;
use crate::adapter::sqlite::ProjectDatabase;
use crate::adapter::sqlite_repos::{SqliteSessionRepository, SqliteWorktreeRepository};
use crate::repository::{NewSession, NewWorktree, WorktreeRecord, WorktreeRepository};
use std::process::Command as StdCommand;
const GIT_STATE_VARS: &[&str] = &[
"GIT_DIR",
"GIT_WORK_TREE",
"GIT_INDEX_FILE",
"GIT_OBJECT_DIRECTORY",
"GIT_ALTERNATE_OBJECT_DIRECTORIES",
"GIT_COMMON_DIR",
"GIT_NAMESPACE",
"GIT_CEILING_DIRECTORIES",
"GIT_AUTHOR_NAME",
"GIT_AUTHOR_EMAIL",
"GIT_AUTHOR_DATE",
"GIT_COMMITTER_NAME",
"GIT_COMMITTER_EMAIL",
"GIT_COMMITTER_DATE",
"GIT_CONFIG_GLOBAL",
"GIT_CONFIG_SYSTEM",
];
fn git_fixture(repo_root: &Path, args: &[&str]) {
let mut cmd = StdCommand::new("git");
cmd.args(args).current_dir(repo_root);
for v in GIT_STATE_VARS {
cmd.env_remove(v);
}
let out = cmd.output().unwrap_or_else(|e| panic!("git {args:?}: {e}"));
assert!(
out.status.success(),
"git {args:?} failed: {}",
String::from_utf8_lossy(&out.stderr)
);
}
fn init_repo() -> (tempfile::TempDir, PathBuf) {
let dir = tempfile::tempdir().expect("tempdir");
let root = dir.path().join("repo");
std::fs::create_dir_all(&root).unwrap();
git_fixture(&root, &["init", "-b", "main"]);
git_fixture(&root, &["config", "user.email", "test@example.com"]);
git_fixture(&root, &["config", "user.name", "Test"]);
std::fs::write(root.join("README.md"), "# test\n").unwrap();
git_fixture(&root, &["add", "."]);
git_fixture(&root, &["commit", "-m", "init"]);
(dir, root)
}
fn seeded_db(dir: &Path) -> ProjectDatabase {
let mut db = ProjectDatabase::open(dir.join("test.sqlite")).unwrap();
db.migrate().unwrap();
db.connection_mut()
.execute(
"INSERT INTO projects (id, name, task_prefix, repository_root, git_common_dir, main_branch, schema_version, created_at, updated_at)
VALUES ('p1', 'proj', 'CTX', '/tmp/r1', '/tmp/g1', 'main', 1, 'now', 'now')",
[],
)
.unwrap();
db
}
struct FailingIdWorktreeRepository;
impl WorktreeRepository for FailingIdWorktreeRepository {
fn upsert(
&self,
_input: &NewWorktree,
_now: &str,
) -> Result<WorktreeRecord, CarryCtxError> {
unreachable!("not used by cleanup test")
}
fn find_by_id(
&self,
_project_id: &str,
_id: &str,
) -> Result<Option<WorktreeRecord>, CarryCtxError> {
Err(CarryCtxError::database_error("injected id lookup failure"))
}
fn find_by_path(
&self,
_project_id: &str,
_path: &str,
) -> Result<Option<WorktreeRecord>, CarryCtxError> {
unreachable!("id lookup must fail before path lookup")
}
fn find_by_task_id(
&self,
_project_id: &str,
_task_id: &str,
) -> Result<Option<WorktreeRecord>, CarryCtxError> {
unreachable!("not used by cleanup test")
}
fn list(&self, _project_id: &str) -> Result<Vec<WorktreeRecord>, CarryCtxError> {
unreachable!("not used by cleanup test")
}
fn unbind_task(
&self,
_id: &str,
_project_id: &str,
_now: &str,
) -> Result<WorktreeRecord, CarryCtxError> {
unreachable!("not used by cleanup test")
}
fn delete(&self, _id: &str, _project_id: &str) -> Result<(), CarryCtxError> {
unreachable!("not used by cleanup test")
}
fn prune_stale(
&self,
_project_id: &str,
_repository_root: &Path,
_actor_agent_id: Option<&str>,
_session_id: Option<&str>,
_now: &str,
) -> Result<Vec<WorktreeRecord>, CarryCtxError> {
unreachable!("not used by cleanup test")
}
}
#[test]
fn blocker_db_string_round_trips_via_domain() {
let cases = vec![
CleanupBlocker::ActiveSession {
session_id: "sess-1".into(),
},
CleanupBlocker::DirtyWorktree,
CleanupBlocker::CurrentWorkingDirectory,
CleanupBlocker::MissingGitMetadata,
CleanupBlocker::WorktreeLocked,
];
for b in cases {
let s = b.to_db_string();
let back = CleanupBlocker::from_db_string(&s).expect("round-trip");
assert_eq!(back, b);
assert_eq!(b.to_string(), s);
}
}
#[test]
fn assess_removable_when_no_blockers() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-clean");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/clean", None)
.unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
let conn = db.connection_mut();
conn.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/clean', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
let session_repo = SqliteSessionRepository::new(conn);
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let assessment = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
Some(&PathBuf::from("/tmp")),
)
.unwrap();
assert!(assessment.removable);
assert!(assessment.blockers.is_empty());
}
#[test]
fn assess_detects_active_session_by_worktree_id() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-active");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/active", None)
.unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
{
let conn = db.connection_mut();
conn.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/active', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
conn.execute(
"INSERT INTO agents (id, project_id, name, provider, role, kind, created_at, updated_at)
VALUES ('agent1', 'p1', 'a1', 'test', NULL, NULL, 'now', 'now')",
[],
)
.unwrap();
let sess = SqliteSessionRepository::new(conn);
sess.create(
&NewSession {
id: "sess1".into(),
project_id: "p1".into(),
agent_id: "agent1".into(),
task_id: None,
worktree_id: Some("wt1".into()),
branch: None,
head: None,
cwd: Some(wt.to_string_lossy().to_string()),
provider: None,
},
"now",
)
.unwrap();
}
let conn = db.connection_mut();
let session_repo = SqliteSessionRepository::new(conn);
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let assessment = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
Some(&PathBuf::from("/tmp")),
)
.unwrap();
assert!(!assessment.removable);
assert!(
assessment
.blockers
.iter()
.any(|b| matches!(b, CleanupBlocker::ActiveSession { .. }))
);
}
#[test]
fn assess_detects_active_session_by_cwd() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-cwd");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/cwd", None)
.unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
{
let conn = db.connection_mut();
conn.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/cwd', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
conn.execute(
"INSERT INTO agents (id, project_id, name, provider, role, kind, created_at, updated_at)
VALUES ('agent1', 'p1', 'a1', 'test', NULL, NULL, 'now', 'now')",
[],
)
.unwrap();
let sess = SqliteSessionRepository::new(conn);
sess.create(
&NewSession {
id: "sess2".into(),
project_id: "p1".into(),
agent_id: "agent1".into(),
task_id: None,
worktree_id: None,
branch: None,
head: None,
cwd: Some(wt.to_string_lossy().to_string()),
provider: None,
},
"now",
)
.unwrap();
}
let conn = db.connection_mut();
let session_repo = SqliteSessionRepository::new(conn);
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let assessment = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
Some(&PathBuf::from("/tmp")),
)
.unwrap();
assert!(
assessment
.blockers
.iter()
.any(|b| matches!(b, CleanupBlocker::ActiveSession { .. }))
);
}
#[test]
fn assess_detects_current_working_directory() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-cwd2");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/cwd2", None)
.unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
db.connection_mut()
.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/cwd2', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
let conn = db.connection_mut();
let session_repo = SqliteSessionRepository::new(conn);
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let assessment = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
Some(&wt),
)
.unwrap();
assert!(
assessment
.blockers
.contains(&CleanupBlocker::CurrentWorkingDirectory)
);
}
#[test]
fn assess_detects_dirty_worktree() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-dirty");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/dirty", None)
.unwrap();
std::fs::write(wt.join("untracked.txt"), "hello\n").unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
db.connection_mut()
.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/dirty', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
let conn = db.connection_mut();
let session_repo = SqliteSessionRepository::new(conn);
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let assessment = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
Some(&PathBuf::from("/tmp")),
)
.unwrap();
assert!(assessment.blockers.contains(&CleanupBlocker::DirtyWorktree));
}
#[test]
fn assess_detects_locked_worktree() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-locked");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/locked", None)
.unwrap();
let mut cmd = StdCommand::new("git");
cmd.args(["worktree", "lock", "--reason", "test lock"])
.arg(&wt)
.current_dir(&repo_root);
for v in GIT_STATE_VARS {
cmd.env_remove(v);
}
let out = cmd.output().expect("lock");
assert!(out.status.success(), "lock failed: {:?}", out);
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
db.connection_mut()
.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/locked', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
let conn = db.connection_mut();
let session_repo = SqliteSessionRepository::new(conn);
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let assessment = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
Some(&PathBuf::from("/tmp")),
)
.unwrap();
assert!(
assessment
.blockers
.contains(&CleanupBlocker::WorktreeLocked)
);
}
#[test]
fn assess_missing_git_metadata_when_id_has_no_row() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-exists");
std::fs::create_dir_all(&wt).unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
let conn = db.connection_mut();
let session_repo = SqliteSessionRepository::new(conn);
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let assessment = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("ghost"),
&wt,
&repo_root,
Some(&PathBuf::from("/tmp")),
)
.unwrap();
assert!(
assessment
.blockers
.contains(&CleanupBlocker::MissingGitMetadata)
);
}
#[test]
fn assess_is_idempotent_across_calls() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-idem");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/idem", None)
.unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
db.connection_mut()
.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/idem', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
std::fs::write(wt.join("dirty.txt"), "x\n").unwrap();
let conn = db.connection_mut();
let session_repo = SqliteSessionRepository::new(conn);
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let a1 = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
Some(&PathBuf::from("/tmp")),
)
.unwrap();
let a2 = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
Some(&PathBuf::from("/tmp")),
)
.unwrap();
assert_eq!(a1, a2);
assert!(a1.blockers.contains(&CleanupBlocker::DirtyWorktree));
}
#[test]
fn execute_is_idempotent_when_path_missing() {
let (_tmp, repo_root) = init_repo();
let missing = repo_root.join("does-not-exist");
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
let conn = db.connection_mut();
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let r1 = execute_worktree_cleanup(
&worktree_repo,
&git_cli,
"p1",
Some("ghost"),
&missing,
&repo_root,
false,
)
.unwrap();
assert!(matches!(r1, ExecuteOutcome::AlreadyRemoved));
let r2 = execute_worktree_cleanup(
&worktree_repo,
&git_cli,
"p1",
Some("ghost"),
&missing,
&repo_root,
false,
)
.unwrap();
assert!(matches!(r2, ExecuteOutcome::AlreadyRemoved));
}
#[test]
fn execute_is_idempotent_when_row_missing() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-ghost");
std::fs::create_dir_all(&wt).unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
let conn = db.connection_mut();
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let outcome = execute_worktree_cleanup(
&worktree_repo,
&git_cli,
"p1",
Some("ghost"),
&wt,
&repo_root,
false,
)
.unwrap();
assert!(matches!(outcome, ExecuteOutcome::AlreadyRemoved));
}
#[test]
fn execute_propagates_worktree_id_lookup_errors() {
let (_tmp, repo_root) = init_repo();
let path = repo_root.join("wt-error");
std::fs::create_dir_all(&path).unwrap();
let outcome = execute_worktree_cleanup(
&FailingIdWorktreeRepository,
&GitCli::new(),
"p1",
Some("wt1"),
&path,
&repo_root,
false,
);
let error = outcome.expect_err("lookup failures must not become AlreadyRemoved");
assert_eq!(error.code, "DATABASE_ERROR");
assert_eq!(error.message, "injected id lookup failure");
}
#[test]
fn execute_removes_live_worktree_and_is_idempotent_after() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-live");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/live", None)
.unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
db.connection_mut()
.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/live', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
let conn = db.connection_mut();
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let r1 = execute_worktree_cleanup(
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
false,
)
.unwrap();
assert!(matches!(r1, ExecuteOutcome::Removed));
assert!(!wt.exists());
let r2 = execute_worktree_cleanup(
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
false,
)
.unwrap();
assert!(matches!(r2, ExecuteOutcome::AlreadyRemoved));
}
#[test]
fn execute_maps_locked_to_blocked_not_failed() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-locked-exec");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/locked-exec", None)
.unwrap();
let mut cmd = StdCommand::new("git");
cmd.args(["worktree", "lock", "--reason", "hold"])
.arg(&wt)
.current_dir(&repo_root);
for v in GIT_STATE_VARS {
cmd.env_remove(v);
}
cmd.output().expect("lock");
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
db.connection_mut()
.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/locked-exec', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
let conn = db.connection_mut();
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let outcome = execute_worktree_cleanup(
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
false,
)
.unwrap();
match outcome {
ExecuteOutcome::Blocked(a) => {
assert!(a.blockers.contains(&CleanupBlocker::WorktreeLocked));
}
other => panic!("expected Blocked, got {other:?}"),
}
}
#[test]
fn execute_maps_dirty_to_blocked_not_failed() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-dirty-exec");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/dirty-exec", None)
.unwrap();
std::fs::write(wt.join("untracked.txt"), "x\n").unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
db.connection_mut()
.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/dirty-exec', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
let conn = db.connection_mut();
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let outcome = execute_worktree_cleanup(
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
false,
)
.unwrap();
match outcome {
ExecuteOutcome::Blocked(a) => {
assert!(a.blockers.contains(&CleanupBlocker::DirtyWorktree));
}
other => panic!("expected Blocked for dirty, got {other:?}"),
}
}
#[test]
fn jj_colocation_blocks_assessment_and_execution_without_removing_worktree() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-jj");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/jj", None)
.unwrap();
std::fs::create_dir(repo_root.join(".jj")).unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
db.connection_mut()
.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/jj', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
let conn = db.connection_mut();
let session_repo = SqliteSessionRepository::new(conn);
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let assessment = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
Some(Path::new("/tmp")),
)
.unwrap();
assert!(assessment.blockers.contains(&CleanupBlocker::JjColocation));
let outcome = execute_worktree_cleanup(
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
true,
)
.unwrap();
assert!(matches!(
outcome,
ExecuteOutcome::Blocked(ref assessment)
if assessment.blockers == vec![CleanupBlocker::JjColocation]
));
assert!(wt.exists());
assert!(
git_cli
.list_worktrees(&repo_root)
.unwrap()
.iter()
.any(|entry| Path::new(&entry.path) == wt)
);
}
#[test]
fn execute_treats_not_a_worktree_as_already_removed() {
let (_tmp, repo_root) = init_repo();
let foreign = repo_root.join("foreign");
std::fs::create_dir_all(&foreign).unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
db.connection_mut()
.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'foreign', 'abc', 'now', 'now')",
rusqlite::params![foreign.to_string_lossy().to_string()],
)
.unwrap();
let conn = db.connection_mut();
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let outcome = execute_worktree_cleanup(
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&foreign,
&repo_root,
false,
)
.unwrap();
assert!(matches!(outcome, ExecuteOutcome::AlreadyRemoved));
}
#[test]
fn assess_does_not_call_remove() {
let (_tmp, repo_root) = init_repo();
let wt = repo_root.join("wt-pure");
GitCli::new()
.create_worktree(&repo_root, &wt, "feat/pure", None)
.unwrap();
std::fs::write(wt.join("dirty.txt"), "x\n").unwrap();
let db_dir = tempfile::tempdir().unwrap();
let mut db = seeded_db(db_dir.path());
db.connection_mut()
.execute(
"INSERT INTO worktrees (id, project_id, task_id, normalized_path, git_common_dir, branch, head, bound_at, updated_at)
VALUES ('wt1', 'p1', NULL, ?1, '', 'feat/pure', 'abc', 'now', 'now')",
rusqlite::params![wt.to_string_lossy().to_string()],
)
.unwrap();
let conn = db.connection_mut();
let session_repo = SqliteSessionRepository::new(conn);
let worktree_repo = SqliteWorktreeRepository::new(conn);
let git_cli = GitCli::new();
let assessment = assess_worktree_cleanup(
&session_repo,
&worktree_repo,
&git_cli,
"p1",
Some("wt1"),
&wt,
&repo_root,
Some(&PathBuf::from("/tmp")),
)
.unwrap();
assert!(wt.exists(), "assess must not remove the worktree");
assert!(assessment.blockers.contains(&CleanupBlocker::DirtyWorktree));
}
}