use std::path::{Path, PathBuf};
use std::sync::Arc;
use async_trait::async_trait;
use tracing::{debug, warn};
use crate::config::OrchestratorConfig;
use crate::vcs::git::commands;
use super::service::{
ConflictPolicy, DirtyState, MergeAttempt, SafetyFact, WorktreeBackend, WorktreeFacts,
WorktreeOpError, WorktreeOpResult,
};
fn utf8_path(path: &Path) -> WorktreeOpResult<&str> {
path.to_str()
.ok_or_else(|| WorktreeOpError::Internal("worktree path is not valid UTF-8".to_string()))
}
pub struct GitWorktreeBackend {
repo_root: PathBuf,
config: Arc<OrchestratorConfig>,
hook_events: Option<tokio::sync::mpsc::Sender<crate::events::ExecutionEvent>>,
}
impl GitWorktreeBackend {
pub fn new(repo_root: PathBuf, config: Arc<OrchestratorConfig>) -> Self {
Self {
repo_root,
config,
hook_events: None,
}
}
pub fn with_hook_events(
mut self,
tx: tokio::sync::mpsc::Sender<crate::events::ExecutionEvent>,
) -> Self {
self.hook_events = Some(tx);
self
}
}
#[async_trait]
impl WorktreeBackend for GitWorktreeBackend {
async fn observe(
&self,
request: crate::worktree_ops::ObservationRequest,
) -> WorktreeOpResult<Vec<WorktreeFacts>> {
let enriched = super::observe_worktrees(&self.repo_root, request)
.await
.map_err(|e| WorktreeOpError::Internal(format!("failed to list worktrees: {e}")))?;
let base_merge_in_progress =
SafetyFact::observed(commands::is_merge_in_progress(&self.repo_root).await);
let mut facts = Vec::with_capacity(enriched.len());
for observation in enriched {
let inspection = observation.info.inspection;
let worktree = observation.info;
let dirty = match commands::has_uncommitted_changes(&worktree.path).await {
Ok((true, _)) => DirtyState::Dirty,
Ok((false, _)) => DirtyState::Clean,
Err(error) => {
debug!(
path = %worktree.path.display(),
"dirty state could not be determined: {error}"
);
DirtyState::Unknown
}
};
facts.push(WorktreeFacts {
identity: crate::parallel::acceptance_state::worktree_identity(&worktree.path),
path: worktree.path,
branch: worktree.branch,
head: worktree.head,
is_main: worktree.is_main,
is_detached: worktree.is_detached,
has_commits_ahead: observation.has_commits_ahead,
conflict_files: worktree
.merge_conflict
.map(|c| c.conflict_files)
.unwrap_or_default(),
dirty,
base_merge_in_progress,
inspection,
});
}
Ok(facts)
}
async fn base_head(&self) -> WorktreeOpResult<String> {
commands::get_current_commit(&self.repo_root)
.await
.map_err(|e| WorktreeOpError::Internal(format!("failed to read base HEAD: {e}")))
}
async fn create(&self, path: &Path, branch: &str, base_commit: &str) -> WorktreeOpResult<()> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).map_err(|e| {
WorktreeOpError::Internal(format!("failed to create workspace directory: {e}"))
})?;
}
let path_str = path.to_str().ok_or_else(|| {
WorktreeOpError::Internal("derived worktree path is not valid UTF-8".to_string())
})?;
commands::worktree_add(&self.repo_root, path_str, branch, base_commit)
.await
.map_err(|e| WorktreeOpError::Internal(format!("failed to create worktree: {e}")))?;
if let Err(error) = commands::run_worktree_setup(&self.repo_root, path).await {
warn!(path = %path.display(), "worktree setup script failed: {error} (non-fatal)");
}
Ok(())
}
async fn teardown(&self, path: &Path) -> WorktreeOpResult<()> {
let path_str = utf8_path(path)?;
commands::run_worktree_teardown(&self.repo_root, path_str)
.await
.map_err(|e| WorktreeOpError::Internal(format!("worktree teardown failed: {e}")))
}
async fn remove_worktree(&self, path: &Path) -> WorktreeOpResult<()> {
let path_str = utf8_path(path)?;
commands::worktree_remove_forced(&self.repo_root, path_str)
.await
.map_err(|e| WorktreeOpError::Internal(format!("failed to remove worktree: {e}")))
}
async fn branch_ref(&self, branch: &str) -> WorktreeOpResult<Option<String>> {
commands::branch_ref_oid(&self.repo_root, branch)
.await
.map_err(|e| WorktreeOpError::Internal(format!("failed to read branch ref: {e}")))
}
async fn delete_branch_if_merged(&self, branch: &str) -> WorktreeOpResult<()> {
commands::branch_delete_if_merged(&self.repo_root, branch)
.await
.map_err(|e| WorktreeOpError::Internal(format!("failed to delete branch: {e}")))
}
async fn delete_branch_at(&self, branch: &str, expected_oid: &str) -> WorktreeOpResult<()> {
commands::branch_delete_at_oid(&self.repo_root, branch, expected_oid)
.await
.map_err(|e| {
WorktreeOpError::Internal(format!(
"failed to delete branch at the confirmed commit: {e}"
))
})
}
async fn merge_into_base(
&self,
branch: &str,
policy: ConflictPolicy,
) -> WorktreeOpResult<MergeAttempt> {
match policy {
ConflictPolicy::AbortOnConflict => {
match commands::merge_branch(&self.repo_root, branch).await {
Ok(()) => Ok(MergeAttempt::Merged),
Err(error) => Err(WorktreeOpError::Internal(error.to_string())),
}
}
ConflictPolicy::PreserveConflict => {
match commands::merge_branch_preserving_conflict(&self.repo_root, branch).await {
Ok(commands::PreservedMergeOutcome::Merged) => Ok(MergeAttempt::Merged),
Ok(commands::PreservedMergeOutcome::Conflict { files }) => {
Ok(MergeAttempt::Conflict { files })
}
Err(error) => Err(WorktreeOpError::Internal(format!(
"base merge failed: {error}"
))),
}
}
}
}
async fn run_on_merged(&self, change_id: &str, worktree_path: &Path) -> WorktreeOpResult<()> {
let hooks_config = self.config.get_hooks();
let hooks = match &self.hook_events {
Some(tx) => crate::hooks::HookRunner::with_event_tx(
hooks_config,
self.repo_root.clone(),
tx.clone(),
),
None => crate::hooks::HookRunner::new(hooks_config, self.repo_root.clone()),
};
let (completed_tasks, total_tasks) =
match crate::openspec::list_changes_native_from(&self.repo_root) {
Ok(changes) => changes
.iter()
.find(|c| c.id == change_id)
.map(|c| (c.completed_tasks, c.total_tasks))
.unwrap_or((0, 0)),
Err(error) => {
warn!("failed to fetch task counts for on_merged hook: {error}");
(0, 0)
}
};
let context = crate::hooks::HookContext::new(0, 0, 0, false)
.with_change(change_id, completed_tasks, total_tasks)
.with_apply_count(0)
.with_parallel_context(&worktree_path.to_string_lossy(), None);
hooks
.run_hook(crate::hooks::HookType::OnMerged, &context)
.await
.map_err(|e| WorktreeOpError::Internal(e.to_string()))
}
async fn change_is_eligible(&self, change_id: &str) -> WorktreeOpResult<()> {
let changes = crate::openspec::list_changes_native_from(&self.repo_root).map_err(|e| {
WorktreeOpError::Internal(format!("failed to enumerate managed changes: {e}"))
})?;
if changes.iter().any(|change| change.id == change_id) {
Ok(())
} else {
Err(WorktreeOpError::NotFound(format!(
"'{change_id}' is not a managed, non-archived change eligible for a worktree"
)))
}
}
}