use std::{
path::{Path, PathBuf},
process::Output,
time::SystemTime,
};
use tracing::instrument;
use zeph_config::{WorktreeBaseRef, WorktreeConfig};
use crate::{
error::WorktreeError,
git_runner::GitRunner,
handle::WorktreeHandle,
sanitize::{canonicalize_root, validate_branch_component},
};
pub struct WorktreeManager<R: GitRunner> {
repo_root: PathBuf,
config: WorktreeConfig,
runner: R,
handles: std::sync::Mutex<Vec<WorktreeHandle>>,
}
impl<R: GitRunner> WorktreeManager<R> {
pub async fn new(
repo_root: PathBuf,
config: WorktreeConfig,
runner: R,
) -> Result<Self, WorktreeError> {
let root = PathBuf::from(&config.root);
let repo = repo_root.clone();
tokio::task::spawn_blocking(move || canonicalize_root(&root, &repo))
.await
.map_err(|e| WorktreeError::Io(std::io::Error::other(e)))??;
Ok(Self {
repo_root,
config,
runner,
handles: std::sync::Mutex::new(Vec::new()),
})
}
#[must_use]
pub fn repo_root(&self) -> &Path {
&self.repo_root
}
#[instrument(name = "worktree.create", skip(self), fields(subagent_id = %subagent_id))]
pub async fn create(&self, subagent_id: &str) -> Result<WorktreeHandle, WorktreeError> {
validate_branch_component(subagent_id)?;
let branch_name = format!("{}{}", self.config.branch_prefix, subagent_id);
let root = PathBuf::from(&self.config.root);
let repo = self.repo_root.clone();
let worktree_root = tokio::task::spawn_blocking(move || canonicalize_root(&root, &repo))
.await
.map_err(|e| WorktreeError::Io(std::io::Error::other(e)))??;
let path = worktree_root.join(subagent_id);
if path.exists() {
return Err(WorktreeError::PathExists(path));
}
let (base_ref_resolved, commitish) = if let WorktreeBaseRef::Fresh = &self.config.base_ref {
let branch = self.resolve_default_branch().await?;
self.fetch_origin(&branch).await?;
self.verify_commitish(&format!("origin/{branch}")).await?;
let resolved = format!("origin/{branch}");
(resolved.clone(), resolved)
} else {
self.check_dirty_tree().await;
("HEAD".to_string(), "HEAD".to_string())
};
let path_str = path.to_string_lossy();
self.git_worktree_add(&branch_name, &path_str, &commitish)
.await?;
let handle = WorktreeHandle {
path,
branch_name,
base_ref_resolved,
subagent_id: subagent_id.to_string(),
created_at: SystemTime::now(),
};
self.handles
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(handle.clone());
Ok(handle)
}
#[instrument(name = "worktree.remove", skip(self), fields(branch = %handle.branch_name))]
pub async fn remove(
&self,
handle: &WorktreeHandle,
prune_branch: bool,
) -> Result<(), WorktreeError> {
let path_str = handle.path.to_string_lossy().to_string();
let out = self
.runner
.run(
&["worktree", "remove", "--force", "--", &path_str],
&self.repo_root,
)
.await?;
check_git_status(&out, "worktree remove")?;
self.handles
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.retain(|h| h.path != handle.path);
if prune_branch {
let branch = &handle.branch_name;
let out = self
.runner
.run(&["branch", "-D", "--", branch], &self.repo_root)
.await?;
check_git_status(&out, "branch -D")?;
}
Ok(())
}
pub fn list(&self) -> Vec<WorktreeHandle> {
self.handles
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
#[instrument(name = "worktree.reconcile", skip(self))]
pub async fn reconcile(&self) -> Result<Vec<WorktreeHandle>, WorktreeError> {
let out = self
.runner
.run(&["worktree", "list", "--porcelain"], &self.repo_root)
.await?;
check_git_status(&out, "worktree list")?;
let output_str = String::from_utf8_lossy(&out.stdout);
let git_worktrees = parse_worktree_list_porcelain(&output_str);
let session_paths: std::collections::HashSet<PathBuf> = self
.handles
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.iter()
.map(|h| h.path.clone())
.collect();
let stale = git_worktrees
.into_iter()
.filter(|h| !session_paths.contains(&h.path))
.filter(|h| h.path != self.repo_root)
.collect();
Ok(stale)
}
}
fn parse_worktree_list_porcelain(output: &str) -> Vec<WorktreeHandle> {
let mut result = Vec::new();
let mut path: Option<PathBuf> = None;
let mut branch: Option<String> = None;
for line in output.lines() {
if let Some(p) = line.strip_prefix("worktree ") {
if let (Some(wt_path), Some(br)) = (path.take(), branch.take()) {
result.push(WorktreeHandle {
path: wt_path,
branch_name: br,
base_ref_resolved: String::new(),
subagent_id: String::new(),
created_at: SystemTime::UNIX_EPOCH,
});
}
path = Some(PathBuf::from(p));
} else if let Some(b) = line.strip_prefix("branch refs/heads/") {
branch = Some(b.to_string());
}
}
if let (Some(wt_path), Some(br)) = (path, branch) {
result.push(WorktreeHandle {
path: wt_path,
branch_name: br,
base_ref_resolved: String::new(),
subagent_id: String::new(),
created_at: SystemTime::UNIX_EPOCH,
});
}
result
}
impl<R: GitRunner> WorktreeManager<R> {
#[instrument(name = "worktree.dirty_check", skip(self))]
async fn check_dirty_tree(&self) {
match self
.runner
.run(&["status", "--porcelain"], &self.repo_root)
.await
{
Ok(out) if !out.stdout.is_empty() => {
tracing::warn!(
"creating a head worktree on a dirty working tree; \
uncommitted changes will NOT be visible in the worktree"
);
}
_ => {}
}
}
#[instrument(name = "worktree.resolve_branch", skip(self))]
async fn resolve_default_branch(&self) -> Result<String, WorktreeError> {
if !self.config.default_branch.is_empty() {
return Ok(self.config.default_branch.clone());
}
let out = self
.runner
.run(
&["symbolic-ref", "refs/remotes/origin/HEAD"],
&self.repo_root,
)
.await?;
if out.status.success() {
let raw = String::from_utf8_lossy(&out.stdout);
let trimmed = raw.trim();
if let Some(branch) = trimmed.strip_prefix("refs/remotes/origin/") {
return Ok(branch.to_string());
}
}
Err(WorktreeError::BaseRefUnresolved {
attempted: "symbolic-ref refs/remotes/origin/HEAD".to_string(),
})
}
#[instrument(name = "worktree.fetch", skip(self), fields(branch = %branch))]
async fn fetch_origin(&self, branch: &str) -> Result<(), WorktreeError> {
let out = self
.runner
.run(&["fetch", "origin", "--", branch], &self.repo_root)
.await?;
check_git_status(&out, "fetch")?;
Ok(())
}
#[instrument(name = "worktree.verify_commitish", skip(self), err)]
async fn verify_commitish(&self, commitish: &str) -> Result<(), WorktreeError> {
let out = self
.runner
.run(&["rev-parse", "--verify", "--", commitish], &self.repo_root)
.await?;
check_git_status(&out, &format!("rev-parse --verify {commitish}"))?;
Ok(())
}
#[instrument(name = "worktree.git_worktree_add", skip(self), err)]
async fn git_worktree_add(
&self,
branch: &str,
path: &str,
commitish: &str,
) -> Result<(), WorktreeError> {
let out = self
.runner
.run(
&["worktree", "add", "-b", branch, "--", path, commitish],
&self.repo_root,
)
.await?;
check_git_status(&out, "worktree add")?;
Ok(())
}
}
fn check_git_status(out: &Output, op: &str) -> Result<(), WorktreeError> {
if !out.status.success() {
let stderr = String::from_utf8_lossy(&out.stderr).to_string();
tracing::debug!(op, %stderr, "git command failed");
return Err(WorktreeError::GitCommand {
op: op.to_string(),
stderr,
});
}
Ok(())
}
#[instrument(name = "worktree.probe_capabilities", skip(runner), err)]
pub async fn probe_capabilities<R: GitRunner>(
runner: &R,
repo_root: &Path,
) -> Result<(), WorktreeError> {
let out = runner.run(&["--version"], repo_root).await?;
if !out.status.success() {
return Err(WorktreeError::GitCommand {
op: "--version".to_string(),
stderr: String::from_utf8_lossy(&out.stderr).to_string(),
});
}
let version_output = String::from_utf8_lossy(&out.stdout);
if let Some(version) = parse_git_version(&version_output)
&& version < (2, 5)
{
return Err(WorktreeError::GitCommand {
op: "--version".to_string(),
stderr: format!(
"git \u{2265} 2.5 is required for worktree support (found: {}.{}). \
Upgrade git or set `worktree.enabled = false`.",
version.0, version.1
),
});
}
let out = runner
.run(&["rev-parse", "--is-inside-work-tree"], repo_root)
.await?;
if !out.status.success() {
return Err(WorktreeError::NotAGitRepo);
}
Ok(())
}
fn parse_git_version(output: &str) -> Option<(u32, u32)> {
let version_str = output.trim().strip_prefix("git version ")?;
let mut parts = version_str.split('.');
let major: u32 = parts.next()?.parse().ok()?;
let minor: u32 = parts.next()?.parse().ok()?;
Some((major, minor))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::git_runner::FakeGitRunner;
use std::assert_matches;
use std::sync::Arc;
use zeph_config::WorktreeConfig;
fn test_config() -> WorktreeConfig {
WorktreeConfig {
enabled: true,
root: "worktrees".to_string(),
branch_prefix: "agent/".to_string(),
..WorktreeConfig::default()
}
}
fn make_repo() -> tempfile::TempDir {
let dir = tempfile::tempdir().unwrap();
std::fs::create_dir_all(dir.path().join(".git")).unwrap();
dir
}
async fn make_manager(
dir: &tempfile::TempDir,
runner: FakeGitRunner,
) -> WorktreeManager<FakeGitRunner> {
WorktreeManager::new(dir.path().to_path_buf(), test_config(), runner)
.await
.unwrap()
}
#[tokio::test]
async fn probe_succeeds_on_valid_git() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_ok(b"git version 2.43.0\n" as &[u8]);
runner.push_ok(b"true\n" as &[u8]);
probe_capabilities(&runner, dir.path()).await.unwrap();
let calls = runner.calls.lock().unwrap();
assert!(calls[0].0.contains(&"--version".to_string()));
assert!(calls[1].0.contains(&"--is-inside-work-tree".to_string()));
}
#[tokio::test]
async fn probe_rejects_old_git() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_ok(b"git version 2.4.0\n" as &[u8]);
let err = probe_capabilities(&runner, dir.path()).await.unwrap_err();
assert_matches!(err, WorktreeError::GitCommand { .. });
}
#[tokio::test]
async fn probe_rejects_non_repo() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_ok(b"git version 2.44.0\n" as &[u8]);
runner.push_err(b"not a git repo\n" as &[u8]);
let err = probe_capabilities(&runner, dir.path()).await.unwrap_err();
assert_matches!(err, WorktreeError::NotAGitRepo);
}
#[tokio::test]
async fn create_head_mode_passes_double_dash() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_ok(b"" as &[u8]);
runner.push_ok(b"" as &[u8]);
let mgr = make_manager(&dir, runner).await;
let result = mgr.create("agent-42").await;
match result {
Ok(_) | Err(WorktreeError::GitCommand { .. }) => {}
Err(e) => panic!("unexpected error: {e}"),
}
}
#[tokio::test]
async fn create_rejects_invalid_branch_component() {
let dir = make_repo();
let runner = FakeGitRunner::new();
let mgr = make_manager(&dir, runner).await;
let err = mgr.create("../escape").await.unwrap_err();
assert_matches!(err, WorktreeError::InvalidBranchName(_));
}
#[tokio::test]
async fn create_rejects_leading_dash() {
let dir = make_repo();
let runner = FakeGitRunner::new();
let mgr = make_manager(&dir, runner).await;
let err = mgr.create("-bad-id").await.unwrap_err();
assert_matches!(err, WorktreeError::InvalidBranchName(_));
}
#[tokio::test]
async fn create_fresh_resolves_default_branch_from_config() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_ok(b"" as &[u8]);
runner.push_ok(b"deadbeef\n" as &[u8]);
runner.push_ok(b"" as &[u8]);
let config = WorktreeConfig {
enabled: true,
base_ref: zeph_config::WorktreeBaseRef::Fresh,
default_branch: "main".to_string(),
root: "worktrees".to_string(),
branch_prefix: "agent/".to_string(),
..WorktreeConfig::default()
};
let mgr = WorktreeManager::new(dir.path().to_path_buf(), config, runner)
.await
.unwrap();
let result = mgr.create("agent-fresh").await;
match result {
Ok(_) | Err(WorktreeError::GitCommand { .. }) => {}
Err(e) => panic!("unexpected error: {e}"),
}
}
#[tokio::test]
async fn create_fresh_fails_when_fetch_fails() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_err(b"network error\n" as &[u8]);
let config = WorktreeConfig {
enabled: true,
base_ref: zeph_config::WorktreeBaseRef::Fresh,
default_branch: "main".to_string(),
root: "worktrees".to_string(),
branch_prefix: "agent/".to_string(),
..WorktreeConfig::default()
};
let mgr = WorktreeManager::new(dir.path().to_path_buf(), config, runner)
.await
.unwrap();
let err = mgr.create("agent-fresh").await.unwrap_err();
assert_matches!(err, WorktreeError::GitCommand { .. });
}
#[tokio::test]
async fn create_fresh_fails_when_symbolic_ref_unset() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_err(b"symbolic-ref: not a ref\n" as &[u8]);
let config = WorktreeConfig {
enabled: true,
base_ref: zeph_config::WorktreeBaseRef::Fresh,
default_branch: String::new(), root: "worktrees".to_string(),
branch_prefix: "agent/".to_string(),
..WorktreeConfig::default()
};
let mgr = WorktreeManager::new(dir.path().to_path_buf(), config, runner)
.await
.unwrap();
let err = mgr.create("agent-fresh").await.unwrap_err();
assert_matches!(err, WorktreeError::BaseRefUnresolved { .. });
}
#[tokio::test]
async fn remove_without_branch_prune() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_ok(b"" as &[u8]);
let mgr = make_manager(&dir, runner).await;
let handle = WorktreeHandle {
path: dir.path().join("worktrees/agent-99"),
branch_name: "agent/agent-99".to_string(),
base_ref_resolved: "HEAD".to_string(),
subagent_id: "agent-99".to_string(),
created_at: SystemTime::now(),
};
mgr.remove(&handle, false).await.unwrap();
}
#[tokio::test]
async fn remove_with_branch_prune_issues_two_git_calls() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_ok(b"" as &[u8]);
runner.push_ok(b"" as &[u8]);
let mgr = make_manager(&dir, runner).await;
let handle = WorktreeHandle {
path: dir.path().join("worktrees/agent-99"),
branch_name: "agent/agent-99".to_string(),
base_ref_resolved: "HEAD".to_string(),
subagent_id: "agent-99".to_string(),
created_at: SystemTime::now(),
};
mgr.remove(&handle, true).await.unwrap();
}
#[tokio::test]
async fn remove_drops_handle_even_when_branch_prune_fails() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_ok(b"" as &[u8]);
runner.push_err(b"error: branch 'agent/agent-99' not fully merged\n" as &[u8]);
let mgr = make_manager(&dir, runner).await;
let handle = WorktreeHandle {
path: dir.path().join("worktrees/agent-99"),
branch_name: "agent/agent-99".to_string(),
base_ref_resolved: "HEAD".to_string(),
subagent_id: "agent-99".to_string(),
created_at: SystemTime::now(),
};
mgr.handles.lock().unwrap().push(handle.clone());
assert_eq!(mgr.list().len(), 1, "precondition: handle is tracked");
let err = mgr.remove(&handle, true).await.unwrap_err();
assert_matches!(
err,
WorktreeError::GitCommand { ref op, .. } if op == "branch -D"
);
assert!(
mgr.list().is_empty(),
"handle must be dropped once `worktree remove` succeeded, \
independent of the branch -D outcome"
);
}
#[tokio::test]
async fn reconcile_parses_porcelain_output() {
let dir = make_repo();
let runner = FakeGitRunner::new();
let porcelain = format!(
"worktree {0}\nHEAD abc123\nbranch refs/heads/main\n\nworktree {0}/worktrees/agent-1\nHEAD def456\nbranch refs/heads/agent/agent-1\n\n",
dir.path().display()
);
runner.push_ok(porcelain.into_bytes());
let mgr = make_manager(&dir, runner).await;
let stale = mgr.reconcile().await.unwrap();
assert_eq!(stale.len(), 1);
assert_eq!(stale[0].branch_name, "agent/agent-1");
}
#[test]
fn parse_version_standard() {
assert_eq!(parse_git_version("git version 2.43.0"), Some((2, 43)));
}
#[test]
fn parse_version_old() {
assert_eq!(parse_git_version("git version 2.4.1"), Some((2, 4)));
}
#[test]
fn parse_version_invalid() {
assert_eq!(parse_git_version("not git output"), None);
}
#[tokio::test]
async fn remove_uses_double_dash_separator() {
let dir = make_repo();
let runner = Arc::new(FakeGitRunner::new());
runner.push_ok(b"" as &[u8]);
let mgr =
WorktreeManager::new(dir.path().to_path_buf(), test_config(), Arc::clone(&runner))
.await
.unwrap();
let handle = WorktreeHandle {
path: dir.path().join("worktrees/x"),
branch_name: "agent/x".to_string(),
base_ref_resolved: "HEAD".to_string(),
subagent_id: "x".to_string(),
created_at: SystemTime::now(),
};
let _ = mgr.remove(&handle, false).await;
let calls = runner.calls.lock().unwrap();
let has_sep = calls[0].0.iter().any(|a| a == "--");
assert!(
has_sep,
"expected '--' separator in git args: {:?}",
calls[0].0
);
}
#[tokio::test]
async fn create_head_mode_proceeds_on_dirty_tree() {
let dir = make_repo();
let runner = FakeGitRunner::new();
runner.push_ok(b" M some-file.txt\n" as &[u8]);
runner.push_err(b"fake error\n" as &[u8]);
let mgr = make_manager(&dir, runner).await;
let result = mgr.create("dirty-agent").await;
assert!(
matches!(result, Err(WorktreeError::GitCommand { .. })),
"expected GitCommand error from fake runner, not an early abort: {result:?}"
);
}
}