use crate::error::{BallError, Result};
use crate::policy::ClaimPolicy;
use crate::store::{task_lock, Store};
use crate::task::{self, Status};
use crate::{claim_sync, git};
use std::{fs, path::PathBuf};
pub use crate::worktree_teardown::{cleanup_orphans, drop_no_worktree, drop_worktree};
pub(crate) fn with_task_lock<T>(
store: &Store,
id: &str,
f: impl FnOnce() -> Result<T>,
) -> Result<T> {
let _guard = task_lock(store, id)?;
f()
}
pub(crate) fn claim_file_path(store: &Store, id: &str) -> PathBuf {
store.claims_dir().join(id)
}
fn write_claim_file(store: &Store, id: &str, worker: &str) -> Result<()> {
fs::create_dir_all(store.claims_dir())?;
let content = format!(
"worker={}\npid={}\nclaimed_at={}\n",
worker,
std::process::id(),
chrono::Utc::now().format("%Y-%m-%dT%H:%M:%SZ")
);
fs::write(claim_file_path(store, id), content)?;
Ok(())
}
pub(crate) fn worktree_path(store: &Store, id: &str) -> Result<PathBuf> {
task::validate_id(id)?;
Ok(store.worktrees_root()?.join(id))
}
pub fn create_worktree(
store: &Store,
id: &str,
identity: &str,
policy: ClaimPolicy,
) -> Result<PathBuf> {
if !store.task_exists(id) {
return Err(BallError::TaskNotFound(id.to_string()));
}
with_task_lock(store, id, || {
let mut task = store.load_task(id)?;
if task.status != Status::Open {
return Err(BallError::NotClaimable(format!(
"{} (status = {})",
id,
task.status.as_str()
)));
}
if task.claimed_by.is_some() {
return Err(BallError::AlreadyClaimed(id.to_string()));
}
let all = store.all_tasks()?;
if crate::ready::is_dep_blocked(&all, &task) {
return Err(BallError::DepsUnmet(id.to_string()));
}
let wt_path = worktree_path(store, id)?;
if wt_path.exists() {
return Err(BallError::WorktreeExists(wt_path));
}
if claim_file_path(store, id).exists() {
return Err(BallError::AlreadyClaimed(id.to_string()));
}
let branch = format!("work/{id}");
task.status = Status::InProgress;
task.claimed_by = Some(identity.to_string());
task.branch = Some(branch.clone());
task.touch();
store.save_task(&task)?;
store.commit_task(id, &format!("balls: claim {} - {}", id, task.title))?;
if policy.require_remote && !store.stealth {
sync_or_rollback(store, id, identity)?;
}
if let Some(parent) = wt_path.parent() {
fs::create_dir_all(parent)?;
}
git::git_worktree_add(&store.root, &wt_path, &branch).inspect_err(|_| {
let _ = rollback_claim(store, id);
})?;
link_shared_state(store, &wt_path)?;
write_claim_file(store, id, identity)?;
Ok(wt_path.clone())
})
}
fn sync_or_rollback(store: &Store, id: &str, identity: &str) -> Result<()> {
match claim_sync::push_claim(store, id, identity) {
Ok(claim_sync::SyncedClaimResult::Pushed) => Ok(()),
Ok(claim_sync::SyncedClaimResult::Lost { winner }) => {
Err(BallError::AlreadyClaimed(format!("{id} (won by {winner})")))
}
Err(e) => {
let _ = rollback_claim(store, id);
Err(e)
}
}
}
fn link_shared_state(store: &Store, wt_path: &std::path::Path) -> Result<()> {
let wt_balls = wt_path.join(".balls");
fs::create_dir_all(&wt_balls)?;
link_state_path(store.local_dir(), &wt_balls.join("local"))?;
if !store.stealth {
link_state_path(store.state_worktree_dir(), &wt_balls.join("worktree"))?;
link_state_path(PathBuf::from("worktree/.balls/tasks"), &wt_balls.join("tasks"))?;
}
Ok(())
}
fn link_state_path(src: PathBuf, dst: &std::path::Path) -> Result<()> {
use std::os::unix::fs::symlink;
if dst.is_symlink() {
return Ok(());
}
if dst.exists() {
return Err(BallError::Other(format!(
"unexpected non-symlink at {}; refusing to link state into worktree",
dst.display()
)));
}
symlink(src, dst)?;
Ok(())
}
fn rollback_claim(store: &Store, id: &str) -> Result<()> {
if let Ok(mut t) = store.load_task(id) {
t.status = Status::Open;
t.claimed_by = None;
t.branch = None;
t.touch();
store.save_task(&t)?;
let _ = store.commit_task(id, &format!("balls: rollback claim {id}"));
}
let _ = fs::remove_file(claim_file_path(store, id));
Ok(())
}
pub fn claim_no_worktree(
store: &Store,
id: &str,
identity: &str,
policy: ClaimPolicy,
) -> Result<()> {
if !store.task_exists(id) {
return Err(BallError::TaskNotFound(id.to_string()));
}
with_task_lock(store, id, || {
let mut task = store.load_task(id)?;
if task.status != Status::Open {
return Err(BallError::NotClaimable(format!("{} (status = {})", id, task.status.as_str())));
}
if task.claimed_by.is_some() {
return Err(BallError::AlreadyClaimed(id.to_string()));
}
let all = store.all_tasks()?;
if crate::ready::is_dep_blocked(&all, &task) {
return Err(BallError::DepsUnmet(id.to_string()));
}
task.status = Status::InProgress;
task.claimed_by = Some(identity.to_string());
task.touch();
store.save_task(&task)?;
store.commit_task(id, &format!("balls: claim {} - {}", id, task.title))?;
if policy.require_remote && !store.stealth {
sync_or_rollback(store, id, identity)?;
}
write_claim_file(store, id, identity)?;
Ok(())
})
}
#[cfg(test)]
#[path = "worktree_tests.rs"]
mod tests;