use std::collections::{HashMap, HashSet};
use std::path::{Component, Path, PathBuf};
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use utoipa::ToSchema;
use crate::worktree_ops::service::{
classify_delete_eligibility, classify_merge_eligibility, ConflictPolicy, DeleteOptions,
ExpectedTarget, WorktreeFacts, WorktreeOpError, WorktreeService, RECOVERY_LOCAL_OR_TUI,
};
use super::dto::{new_hex_id, ErrorCode};
use super::executor::{CommandFailure, ExecutionSummary};
pub const WORKTREE_OPERATIONS: [&str; 5] = ["list", "detail", "create", "delete", "merge"];
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct WorktreeKey {
pub path: PathBuf,
pub identity: String,
}
impl WorktreeKey {
pub fn from_facts(facts: &WorktreeFacts) -> Self {
Self {
path: facts.path.clone(),
identity: facts.identity.clone(),
}
}
}
#[derive(Default)]
struct RegistryInner {
by_id: HashMap<String, WorktreeKey>,
by_key: HashMap<WorktreeKey, String>,
retired: HashSet<String>,
}
#[derive(Default)]
pub struct WorktreeRegistry {
inner: Mutex<RegistryInner>,
}
impl WorktreeRegistry {
pub fn new() -> Self {
Self::default()
}
pub fn sync(&self, observed: &[WorktreeKey]) -> Vec<String> {
let mut inner = self.lock();
let present: HashSet<&WorktreeKey> = observed.iter().collect();
let disappeared: Vec<WorktreeKey> = inner
.by_key
.keys()
.filter(|key| !present.contains(*key))
.cloned()
.collect();
for key in disappeared {
if let Some(id) = inner.by_key.remove(&key) {
inner.by_id.remove(&id);
inner.retired.insert(id);
}
}
observed
.iter()
.map(|key| match inner.by_key.get(key) {
Some(id) => id.clone(),
None => {
let id = Self::allocate(&inner.retired);
inner.by_key.insert(key.clone(), id.clone());
inner.by_id.insert(id.clone(), key.clone());
id
}
})
.collect()
}
pub fn resolve(&self, worktree_id: &str) -> Option<WorktreeKey> {
self.lock().by_id.get(worktree_id).cloned()
}
pub fn retire(&self, worktree_id: &str) {
let mut inner = self.lock();
if let Some(key) = inner.by_id.remove(worktree_id) {
inner.by_key.remove(&key);
}
inner.retired.insert(worktree_id.to_string());
}
pub fn is_retired(&self, worktree_id: &str) -> bool {
self.lock().retired.contains(worktree_id)
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn len(&self) -> usize {
self.lock().by_id.len()
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
fn allocate(retired: &HashSet<String>) -> String {
loop {
let id = new_hex_id();
if !retired.contains(&id) {
return id;
}
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, RegistryInner> {
self.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
}
const FNV_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
pub fn repository_correlation_id(repo_root: &Path) -> String {
let identity = crate::parallel::acceptance_state::repository_identity(repo_root);
let mut hash = FNV_OFFSET_BASIS;
for byte in identity.as_bytes() {
hash ^= *byte as u64;
hash = hash.wrapping_mul(FNV_PRIME);
}
format!("{hash:016x}")
}
pub fn repository_relative_display(repo_root: &Path, target: &Path) -> String {
let root: Vec<Component<'_>> = repo_root
.components()
.filter(|c| !matches!(c, Component::CurDir))
.collect();
let path: Vec<Component<'_>> = target
.components()
.filter(|c| !matches!(c, Component::CurDir))
.collect();
let shared = root
.iter()
.zip(path.iter())
.take_while(|(a, b)| a == b)
.count();
let mut parts: Vec<String> =
std::iter::repeat_n("..".to_string(), root.len() - shared).collect();
parts.extend(
path[shared..]
.iter()
.map(|c| c.as_os_str().to_string_lossy().to_string()),
);
if parts.is_empty() {
".".to_string()
} else {
parts.join("/")
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct WorktreeConflict {
pub files: Vec<String>,
pub recovery: String,
}
impl WorktreeConflict {
pub fn new(files: Vec<String>) -> Self {
Self {
files,
recovery: RECOVERY_LOCAL_OR_TUI.to_string(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct WorktreeEligibility {
pub deletable: bool,
pub mergeable: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub delete_blocked_reason: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub merge_blocked_reason: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
pub struct WorktreeResource {
pub worktree_id: String,
pub repository_id: String,
pub path: String,
pub branch: String,
pub head: String,
pub is_main: bool,
pub is_detached: bool,
pub dirty: Option<bool>,
pub has_commits_ahead: bool,
#[serde(default)]
pub inspection: crate::worktree_ops::InspectionState,
#[serde(skip_serializing_if = "Option::is_none")]
pub conflict: Option<WorktreeConflict>,
pub operations: WorktreeEligibility,
}
impl WorktreeResource {
pub fn project(
worktree_id: String,
repository_id: String,
repo_root: &Path,
facts: &WorktreeFacts,
) -> Self {
let delete = classify_delete_eligibility(facts, DeleteOptions::fail_closed());
let merge = classify_merge_eligibility(facts, ConflictPolicy::PreserveConflict);
Self {
worktree_id,
repository_id,
path: repository_relative_display(repo_root, &facts.path),
branch: facts.branch.clone(),
head: facts.head.clone(),
is_main: facts.is_main,
is_detached: facts.is_detached,
dirty: facts.dirty.as_option(),
has_commits_ahead: facts.has_commits_ahead.is_known_yes(),
inspection: facts.inspection,
conflict: if facts.conflict_files.is_empty() {
None
} else {
Some(WorktreeConflict::new(facts.conflict_files.clone()))
},
operations: WorktreeEligibility {
deletable: delete.is_ok(),
mergeable: merge.is_ok(),
delete_blocked_reason: delete.err().map(|e| e.to_string()),
merge_blocked_reason: merge.err().map(|e| e.to_string()),
},
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct WorktreesResponse {
pub instance_id: String,
pub state_revision: u64,
pub repository_id: String,
pub worktrees: Vec<WorktreeResource>,
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct WorktreeResponse {
pub instance_id: String,
pub state_revision: u64,
pub worktree: WorktreeResource,
}
#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
pub struct WorktreeCapabilities {
pub operations: Vec<String>,
pub merge_conflict_recovery: String,
pub merge_conflict_preserves_state: bool,
pub delete_requires_teardown: bool,
}
impl Default for WorktreeCapabilities {
fn default() -> Self {
Self {
operations: WORKTREE_OPERATIONS.iter().map(|o| o.to_string()).collect(),
merge_conflict_recovery: RECOVERY_LOCAL_OR_TUI.to_string(),
merge_conflict_preserves_state: true,
delete_requires_teardown: true,
}
}
}
#[async_trait]
pub trait WorktreeOperations: Send + Sync {
async fn list(&self) -> Result<WorktreeListing, CommandFailure>;
async fn create(&self, change_id: &str) -> Result<ExecutionSummary, CommandFailure>;
async fn delete(&self, worktree_id: &str) -> Result<ExecutionSummary, CommandFailure>;
async fn merge(&self, worktree_id: &str) -> Result<ExecutionSummary, CommandFailure>;
}
#[derive(Debug, Clone)]
pub struct WorktreeListing {
pub repository_id: String,
pub worktrees: Vec<WorktreeResource>,
}
pub struct UnboundWorktreeOperations;
fn unbound() -> CommandFailure {
CommandFailure::new(
ErrorCode::LifecycleConflict,
"this instance has no worktree runtime bound yet",
)
}
#[async_trait]
impl WorktreeOperations for UnboundWorktreeOperations {
async fn list(&self) -> Result<WorktreeListing, CommandFailure> {
Err(unbound())
}
async fn create(&self, _change_id: &str) -> Result<ExecutionSummary, CommandFailure> {
Err(unbound())
}
async fn delete(&self, _worktree_id: &str) -> Result<ExecutionSummary, CommandFailure> {
Err(unbound())
}
async fn merge(&self, _worktree_id: &str) -> Result<ExecutionSummary, CommandFailure> {
Err(unbound())
}
}
pub fn map_worktree_error(error: &WorktreeOpError) -> CommandFailure {
let code = match error {
WorktreeOpError::NotFound(_) => ErrorCode::WorktreeNotFound,
WorktreeOpError::Exists(_) => ErrorCode::WorktreeExists,
WorktreeOpError::Dirty { .. } => ErrorCode::WorktreeDirty,
WorktreeOpError::DirtyUnknown(_) => ErrorCode::WorktreeDirtyUnknown,
WorktreeOpError::CommitsAhead { .. } | WorktreeOpError::Ineligible(_) => {
ErrorCode::TargetIneligible
}
WorktreeOpError::RootBusy(_) => ErrorCode::RootBusy,
WorktreeOpError::MergeConflict { .. } => ErrorCode::MergeConflict,
WorktreeOpError::Internal(_) => ErrorCode::InternalError,
};
CommandFailure::new(code, error.to_string())
}
pub struct RemoteWorktreeOperations {
service: Arc<WorktreeService>,
registry: Arc<WorktreeRegistry>,
repo_root: PathBuf,
}
impl RemoteWorktreeOperations {
pub fn new(
service: Arc<WorktreeService>,
registry: Arc<WorktreeRegistry>,
repo_root: PathBuf,
) -> Self {
Self {
service,
registry,
repo_root: std::fs::canonicalize(&repo_root).unwrap_or(repo_root),
}
}
async fn observe(&self) -> Result<(WorktreeListing, Vec<WorktreeFacts>), CommandFailure> {
let facts = self
.service
.observe()
.await
.map_err(|error| map_worktree_error(&error))?;
let keys: Vec<WorktreeKey> = facts.iter().map(WorktreeKey::from_facts).collect();
let ids = self.registry.sync(&keys);
let repository_id = repository_correlation_id(&self.repo_root);
let worktrees = ids
.iter()
.zip(facts.iter())
.map(|(id, facts)| {
WorktreeResource::project(id.clone(), repository_id.clone(), &self.repo_root, facts)
})
.collect();
Ok((
WorktreeListing {
repository_id,
worktrees,
},
facts,
))
}
fn resolve(&self, worktree_id: &str) -> Result<WorktreeKey, CommandFailure> {
self.registry.resolve(worktree_id).ok_or_else(|| {
let detail = if self.registry.is_retired(worktree_id) {
"was retired when its worktree disappeared"
} else {
"is not known to this instance"
};
CommandFailure::new(
ErrorCode::WorktreeNotFound,
format!("worktree '{worktree_id}' {detail}"),
)
})
}
}
#[async_trait]
impl WorktreeOperations for RemoteWorktreeOperations {
async fn list(&self) -> Result<WorktreeListing, CommandFailure> {
Ok(self.observe().await?.0)
}
async fn create(&self, change_id: &str) -> Result<ExecutionSummary, CommandFailure> {
let created = self
.service
.create_change_worktree(change_id)
.await
.map_err(|error| map_worktree_error(&error))?;
let (listing, _) = self.observe().await?;
let created_path = repository_relative_display(&self.repo_root, &created.path);
let worktree_id = listing
.worktrees
.iter()
.find(|resource| resource.path == created_path)
.map(|resource| resource.worktree_id.clone())
.ok_or_else(|| {
CommandFailure::new(
ErrorCode::InternalError,
"the created worktree could not be observed for identity allocation",
)
})?;
Ok(ExecutionSummary::changed(format!(
"worktree for change '{change_id}' created as worktree_id {worktree_id}"
)))
}
async fn delete(&self, worktree_id: &str) -> Result<ExecutionSummary, CommandFailure> {
let key = self.resolve(worktree_id)?;
let outcome = self
.service
.delete_worktree(
&key.path,
&ExpectedTarget::unchecked().with_identity(key.identity.clone()),
DeleteOptions::fail_closed(),
)
.await
.map_err(|error| map_worktree_error(&error))?;
self.registry.retire(worktree_id);
let _ = self.observe().await;
Ok(ExecutionSummary::changed(format!(
"{}; worktree_id {worktree_id} retired",
outcome.detail
)))
}
async fn merge(&self, worktree_id: &str) -> Result<ExecutionSummary, CommandFailure> {
let key = self.resolve(worktree_id)?;
let outcome = self
.service
.merge_worktree(&key.path, ConflictPolicy::PreserveConflict)
.await
.map_err(|error| map_worktree_error(&error))?;
let _ = self.observe().await;
Ok(ExecutionSummary::changed(outcome.detail))
}
}