fix(author): bootstrap recovery for dirty orphaned issue worktrees (#860) #861

Merged
sysadmin merged 6 commits from fix/issue-860-dirty-orphan-worktree-recovery into master 2026-07-24 05:52:24 -05:00
9 changed files with 4649 additions and 212 deletions
Showing only changes of commit a54b16676a - Show all commits
File diff suppressed because it is too large Load Diff
+357 -6
View File
@@ -440,13 +440,32 @@ def _session_author_lock_worktree() -> str | None:
Used to derive the author mutation workspace when no explicit Used to derive the author mutation workspace when no explicit
``worktree_path`` or env binding is provided. Never invents a path. ``worktree_path`` or env binding is provided. Never invents a path.
#864: a session pointer whose owner PID is dead and is not this process
must not force workspace binding for other issues rebind is required for
that issue, and a stale dead-owner pointer must not poison unrelated work.
""" """
try: try:
lock = issue_lock_store.read_session_issue_lock() or {} lock = issue_lock_store.read_session_issue_lock() or {}
except Exception: except Exception:
return None return None
path = (lock.get("worktree_path") or "").strip() path = (lock.get("worktree_path") or "").strip()
return path or None if not path:
return None
pid = lock.get("session_pid")
if pid is None:
pid = lock.get("pid")
try:
pid_i = int(pid) if pid is not None else None
except (TypeError, ValueError):
pid_i = None
if (
pid_i is not None
and pid_i != os.getpid()
and not issue_lock_store.is_process_alive(pid_i)
):
return None
return path
def _resolve_preflight_workspace_path(worktree_path: str | None = None) -> str: def _resolve_preflight_workspace_path(worktree_path: str | None = None) -> str:
@@ -2032,6 +2051,7 @@ import issue_lock_adoption # noqa: E402
import issue_lock_recovery # noqa: E402 import issue_lock_recovery # noqa: E402
import issue_lock_renewal # noqa: E402 import issue_lock_renewal # noqa: E402
import dirty_orphan_worktree_recovery # noqa: E402 # #860 dirty orphan recovery import dirty_orphan_worktree_recovery # noqa: E402 # #860 dirty orphan recovery
import dirty_same_claimant_session_rebind # noqa: E402 # #864
import stacked_pr_support # noqa: E402 import stacked_pr_support # noqa: E402
import merge_approval_gate # noqa: E402 import merge_approval_gate # noqa: E402
import review_quarantine # noqa: E402 # #695 contaminated formal-review quarantine import review_quarantine # noqa: E402 # #695 contaminated formal-review quarantine
@@ -4343,6 +4363,7 @@ def gitea_lock_issue(
return result return result
@mcp.tool()
@mcp.tool() @mcp.tool()
def gitea_recover_dirty_orphaned_issue_worktree( def gitea_recover_dirty_orphaned_issue_worktree(
issue_number: int, issue_number: int,
@@ -4616,6 +4637,264 @@ def gitea_recover_dirty_orphaned_issue_worktree(
) )
return result return result
@mcp.tool()
def gitea_rebind_dirty_same_claimant_author_session(
issue_number: int,
branch_name: str,
worktree_path: str,
old_pid: int,
expected_local_head: str,
expected_remote_head: str,
expected_dirty_paths: list[str],
expected_fingerprints: dict,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
dry_run: bool = False,
authorize_reconciler_execute: bool = False,
) -> dict:
"""Rebind a dirty registered issue worktree to this session (#864).
Sanctioned only when every pin agrees: same claimant, dead old_pid matching
the durable lock, matching local/remote heads, exact dirty path set, and
per-path sha256 fingerprints. Preserves every tracked/untracked byte.
Does not sync remote, create recovery worktrees, clean, reset, or move heads.
Role gate:
* author must match the lock claimant identity/profile
* reconciler execute only when ``authorize_reconciler_execute=True``
* reviewer/merger always refuse
``gitea.issue.comment`` (author map entry) is required for mutation; dry_run
still assesses fully but writes nothing. Permission alone is never ownership
proof every pin is re-checked server-side.
Args:
issue_number: Tracking issue number on the durable lock.
branch_name: Exact locked branch name.
worktree_path: Registered dirty worktree path (must be under branches/).
old_pid: Dead owner PID recorded on the lock (must match session_pid/pid).
expected_local_head: Full local HEAD sha the caller observed.
expected_remote_head: Full remote-tracking HEAD sha the caller observed.
expected_dirty_paths: Exact set of dirty relative paths (tracked+untracked).
expected_fingerprints: Map of relative path -> sha256 hex of file bytes.
remote: Known instance 'dadeschools' or 'prgs'.
host/org/repo: Optional target overrides (validated against binding).
dry_run: When true, assess only (no lock/session writes).
authorize_reconciler_execute: Reconciler-only execute gate.
"""
role = _profile_role_kind(get_profile())
role_norm = (role or "").strip().lower()
# Permission: authors need comment; dry_run assess is reachable under read
# for diagnosis, but execute always needs comment. Reconciler execute also
# needs comment when authorized.
if dry_run:
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"dry_run": True,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
else:
blocked = _profile_permission_block(
task_capability_map.required_permission(
"rebind_dirty_same_claimant_author_session"
),
issue_number=issue_number,
remote=remote,
host=host,
org=org,
repo=repo,
org_explicit=org is not None,
repo_explicit=repo is not None,
)
if blocked:
return blocked
if role_norm in {"reviewer", "merger"}:
return {
"success": False,
"dry_run": bool(dry_run),
"outcome": dirty_same_claimant_session_rebind.REFUSED,
"reasons": [
f"role '{role_norm}' cannot rebind dirty same-claimant author "
"sessions (fail closed)"
],
}
if role_norm == "reconciler" and not authorize_reconciler_execute and not dry_run:
return {
"success": False,
"dry_run": False,
"outcome": dirty_same_claimant_session_rebind.REFUSED,
"reasons": [
"reconciler role requires authorize_reconciler_execute=True "
"to execute dirty same-claimant rebind (fail closed)"
],
}
h, o, r = _resolve(remote, host, org, repo)
try:
identity = _authenticated_username(h)
except Exception:
identity = None
profile = get_profile()
profile_name = profile.get("profile_name")
existing = _load_existing_issue_lock(
remote=remote, org=o, repo=r, issue_number=issue_number
)
resolved_wt = os.path.realpath(os.path.abspath((worktree_path or "").strip()))
inv = dirty_same_claimant_session_rebind.collect_dirty_inventory(resolved_wt)
branch_res = subprocess.run(
["git", "-C", resolved_wt, "branch", "--show-current"],
capture_output=True,
text=True,
check=False,
)
current_branch = (branch_res.stdout or "").strip() or None
head_res = subprocess.run(
["git", "-C", resolved_wt, "rev-parse", "HEAD"],
capture_output=True,
text=True,
check=False,
)
local_head = (head_res.stdout or "").strip() if head_res.returncode == 0 else None
# Observe remote-tracking head without network when possible.
remote_head = None
for ref in (
f"refs/remotes/origin/{branch_name}",
f"origin/{branch_name}",
f"refs/remotes/{remote}/{branch_name}",
f"{remote}/{branch_name}",
):
rh = subprocess.run(
["git", "-C", resolved_wt, "rev-parse", "--verify", "--quiet", ref],
capture_output=True,
text=True,
check=False,
)
if rh.returncode == 0 and (rh.stdout or "").strip():
remote_head = (rh.stdout or "").strip()
break
if remote_head is None:
# Fall back to caller's pin only for observation absence — assessment
# still requires pin==observed, so missing observation fails closed.
remote_head = None
# Competing live locks (other issues / other worktrees).
competing_live = []
for entry in issue_lock_store.list_live_locks():
competing_live.append(entry)
# Session pointers that claim this issue lock.
competing_sessions = []
lock_dir = issue_lock_store.default_lock_dir()
lock_path = issue_lock_store.lock_file_path(
remote=remote, org=o, repo=r, issue_number=issue_number, lock_dir=lock_dir
)
try:
for name in os.listdir(lock_dir):
if not name.startswith("session-") or not name.endswith(".json"):
continue
ptr = issue_lock_store.read_lock_file(os.path.join(lock_dir, name))
if not ptr:
continue
ptr_lock = str(ptr.get("lock_file_path") or "").strip()
if not ptr_lock:
continue
try:
same = os.path.realpath(ptr_lock) == os.path.realpath(lock_path)
except OSError:
same = ptr_lock == lock_path
if not same:
continue
try:
sess_pid = int(str(name)[len("session-") : -len(".json")])
except ValueError:
sess_pid = ptr.get("pid")
competing_sessions.append(
{
"pid": sess_pid,
"lock_file_path": ptr_lock,
"live": issue_lock_store.is_process_alive(sess_pid),
}
)
except OSError:
pass
# Best-effort workflow-lease scan: any live lock file whose work_lease is a
# non-author workflow lease on this issue/branch counts as active.
workflow_lease_active = False
for path in issue_lock_store.iter_lock_files(lock_dir):
rec = issue_lock_store.read_lock_file(path)
if not rec:
continue
lease = rec.get("work_lease") if isinstance(rec.get("work_lease"), dict) else {}
op = str(lease.get("operation_type") or "")
if op and op != issue_lock_store.AUTHOR_ISSUE_WORK_LEASE:
if rec.get("issue_number") == issue_number or str(
rec.get("branch_name") or ""
) == branch_name:
if issue_lock_store.is_lease_live(rec):
workflow_lease_active = True
break
repo_root = _canonical_local_git_root()
# permission_allowed reflects profile gate only — never ownership proof.
permission_allowed = True
result = dirty_same_claimant_session_rebind.apply_dirty_same_claimant_session_rebind(
remote=remote,
org=o,
repo=r,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_wt,
claimant_identity=identity,
claimant_profile=profile_name,
old_pid=old_pid,
expected_local_head=expected_local_head,
expected_remote_head=expected_remote_head,
expected_dirty_paths=list(expected_dirty_paths or []),
expected_fingerprints=dict(expected_fingerprints or {}),
existing_lock=existing,
current_identity=identity,
current_profile=profile_name,
role_kind=role_norm or role,
current_pid=os.getpid(),
current_branch=current_branch,
local_head=local_head,
remote_head=remote_head,
dirty_inventory=inv,
competing_live_locks=competing_live,
competing_sessions=competing_sessions,
workflow_lease_active=workflow_lease_active,
authorize_reconciler_execute=bool(authorize_reconciler_execute),
permission_allowed=permission_allowed,
repo_root=repo_root,
dry_run=bool(dry_run),
lock_dir=lock_dir,
)
result["observed"] = {
"local_head": local_head,
"remote_head": remote_head,
"current_branch": current_branch,
"dirty_paths": inv.get("dirty_paths"),
"fingerprints": inv.get("fingerprints"),
"identity": identity,
"profile": profile_name,
"role_kind": role_norm,
}
return result
return result
@mcp.tool() @mcp.tool()
def gitea_assess_work_issue_duplicate( def gitea_assess_work_issue_duplicate(
@@ -11909,6 +12188,7 @@ def gitea_audit_worktree_cleanup(
org: str | None = None, org: str | None = None,
repo: str | None = None, repo: str | None = None,
ttl_hours: float = worktree_cleanup_audit.DEFAULT_TTL_HOURS, ttl_hours: float = worktree_cleanup_audit.DEFAULT_TTL_HOURS,
merged_pr_limit: int = 200,
) -> dict: ) -> dict:
"""Read-only: classify every session-owned worktree under ``branches/`` (#401). """Read-only: classify every session-owned worktree under ``branches/`` (#401).
@@ -11919,17 +12199,26 @@ def gitea_audit_worktree_cleanup(
the active issue-lock branch is read from the local lock file and treated the active issue-lock branch is read from the local lock file and treated
as active work. Deletes nothing and mutates no Gitea state. as active work. Deletes nothing and mutates no Gitea state.
Fails closed if the live open-PR list cannot be fetched: without it, Merged PRs are fetched as well, so an issue worktree can be linked to the
removability cannot be proven, so no candidates are returned. PR that owns its branch (#858). Such a worktree only becomes removable
when that owning PR is unambiguous and merged, the worktree head is
already contained in authoritative master, and nothing else protects it
no open or competing PR, lease, issue lock, live session, dirty file, or
protected/control checkout. Anything unproven keeps it classified as
active issue work.
Fails closed if the live open-PR list, the merged-PR list, or the
control-plane lease state cannot be read: without them removability
cannot be proven, so no candidates are returned.
Args: Args:
remote: Known instance 'dadeschools' or 'prgs'. remote: Known instance 'dadeschools' or 'prgs'.
host: Override the Gitea host. host: Override the Gitea host.
org: Override the owner/organization. org: Override the owner/organization.
repo: Override the repository name. repo: Override the repository name.
ttl_hours: Age (hours) after which a clean issue/conflict-fix ttl_hours: Age (hours) after which a clean conflict-fix worktree
worktree becomes stale-removable (default from becomes stale-removable (default from GITEA_WORKTREE_TTL_HOURS).
GITEA_WORKTREE_TTL_HOURS). merged_pr_limit: Max closed PRs scanned for merged-PR ownership.
Returns: Returns:
dict with per-worktree classifications, counts, removable dict with per-worktree classifications, counts, removable
@@ -11965,22 +12254,84 @@ def gitea_audit_worktree_cleanup(
if (pr.get("head") or {}).get("ref") if (pr.get("head") or {}).get("ref")
} }
# #858: merged PRs are the ownership evidence that lets a landed issue
# worktree stop being reported as active work. Without them the audit can
# never agree with the PR-scoped reconciler, so treat a fetch failure the
# same way an open-PR fetch failure is treated: fail closed.
try:
closed_prs = api_get_all(
f"{repo_api_url(h, o, r)}/pulls?state=closed", auth, limit=merged_pr_limit
)
except Exception as exc:
return {
"success": False,
"performed": False,
"open_pr_state_verified": True,
"merged_pr_state_verified": False,
"reasons": [
"could not fetch merged PRs; worktree ownership unverified "
f"(fail closed): {_redact(str(exc))}"
],
}
merged_prs = [pr for pr in closed_prs if (pr.get("merged") or pr.get("merged_at"))]
pr_index = worktree_cleanup_audit.build_pr_index(list(open_prs) + merged_prs)
# #858: the auditor already accepted lease evidence but nothing ever
# supplied it, so every worktree looked unleased. Removability is now
# reachable for issue worktrees, so authoritative control-plane leases
# must be readable or the audit fails closed.
db, lease_errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"performed": False,
"open_pr_state_verified": True,
"merged_pr_state_verified": True,
"lease_state_verified": False,
"reasons": [
"could not read control-plane leases; worktree protection "
"unverified (fail closed)",
*lease_errs,
],
}
lease_result = lease_lifecycle.list_active_leases(
db, remote=remote, org=o, repo=r, include_non_active=False, limit=500
)
leased_issue_numbers: set[int] = set()
live_session_paths: set[str] = set()
for lease in lease_result.get("leases") or []:
if lease.get("work_kind") == "issue" and lease.get("work_number") is not None:
try:
leased_issue_numbers.add(int(lease["work_number"]))
except (TypeError, ValueError):
pass
if lease.get("worktree_path"):
live_session_paths.add(str(lease["worktree_path"]))
active_issue_branches: set[str] = set() active_issue_branches: set[str] = set()
lock = merged_cleanup_reconcile.read_issue_lock(ISSUE_LOCK_FILE) lock = merged_cleanup_reconcile.read_issue_lock(ISSUE_LOCK_FILE)
if lock and lock.get("branch_name"): if lock and lock.get("branch_name"):
active_issue_branches.add(str(lock["branch_name"]).strip()) active_issue_branches.add(str(lock["branch_name"]).strip())
master_ref = f"{remote}/master" if remote in REMOTES else "origin/master"
report = worktree_cleanup_audit.audit_branches_directory( report = worktree_cleanup_audit.audit_branches_directory(
_canonical_local_git_root(), _canonical_local_git_root(),
open_pr_branches=open_pr_branches, open_pr_branches=open_pr_branches,
active_issue_branches=active_issue_branches, active_issue_branches=active_issue_branches,
now=datetime.now(timezone.utc), now=datetime.now(timezone.utc),
ttl_hours=ttl_hours, ttl_hours=ttl_hours,
pr_index=pr_index,
leased_issue_numbers=leased_issue_numbers,
live_session_paths=live_session_paths,
master_ref=master_ref,
) )
return { return {
"success": True, "success": True,
"performed": False, "performed": False,
"open_pr_state_verified": True, "open_pr_state_verified": True,
"merged_pr_state_verified": True,
"lease_state_verified": True,
"master_ref": master_ref,
"task_mode": "work-issue", "task_mode": "work-issue",
**report, **report,
} }
+5
View File
@@ -17,12 +17,17 @@ SOURCE_LOCK_ISSUE = "gitea_lock_issue"
SOURCE_LOCK_ADOPTION = "gitea_lock_issue_adoption" SOURCE_LOCK_ADOPTION = "gitea_lock_issue_adoption"
SOURCE_OPERATOR_OVERRIDE = "operator_override" SOURCE_OPERATOR_OVERRIDE = "operator_override"
SOURCE_RECOVER_DIRTY_ORPHANED = "gitea_recover_dirty_orphaned_issue_worktree" SOURCE_RECOVER_DIRTY_ORPHANED = "gitea_recover_dirty_orphaned_issue_worktree"
# #864: dirty-preserving same-claimant author-session rebind (dead owner PID).
SOURCE_DIRTY_SAME_CLAIMANT_REBIND = (
"gitea_rebind_dirty_same_claimant_author_session"
)
SANCTIONED_LOCK_SOURCES = frozenset({ SANCTIONED_LOCK_SOURCES = frozenset({
SOURCE_LOCK_ISSUE, SOURCE_LOCK_ISSUE,
SOURCE_LOCK_ADOPTION, SOURCE_LOCK_ADOPTION,
SOURCE_OPERATOR_OVERRIDE, SOURCE_OPERATOR_OVERRIDE,
SOURCE_RECOVER_DIRTY_ORPHANED, SOURCE_RECOVER_DIRTY_ORPHANED,
SOURCE_DIRTY_SAME_CLAIMANT_REBIND,
}) })
_OPERATOR_OVERRIDE_ENV = "GITEA_ISSUE_LOCK_OPERATOR_OVERRIDE" _OPERATOR_OVERRIDE_ENV = "GITEA_ISSUE_LOCK_OPERATOR_OVERRIDE"
+519 -200
View File
@@ -1,30 +1,71 @@
import base64 """Keyed, persistent issue-lock storage (#443) with flock hardening (#438).
import contextlib
Replaces the single global ``/tmp/gitea_issue_lock.json`` slot with per-issue
lock files under ``GITEA_ISSUE_LOCK_DIR`` (default
``~/.cache/gitea-tools/issue-locks``). Each MCP session binds its active lock
via a per-process pointer file so concurrent repos/issues never clobber each
other. Acquisition is serialized per issue with ``fcntl.flock``.
"""
from __future__ import annotations
import errno
import fcntl
import json import json
import os import os
import re import re
import subprocess import tempfile
from contextlib import contextmanager
from datetime import datetime, timedelta, timezone from datetime import datetime, timedelta, timezone
from typing import Any from typing import Any
from audit_event_reconciliation import redact_sensitive_text LOCK_DIR_ENV = "GITEA_ISSUE_LOCK_DIR"
import issue_lock_provenance DEFAULT_LOCK_DIR = os.path.expanduser("~/.cache/gitea-tools/issue-locks")
WORK_LEASE_TTL_HOURS = 4
DEFAULT_LOCK_TTL_HOURS = 24
AUTHOR_ISSUE_WORK_LEASE = "author_issue_work" AUTHOR_ISSUE_WORK_LEASE = "author_issue_work"
_SANCTIONED_ROLES = {"author", "reviewer", "merger", "reconciler", "controller"} _SAFE_SEGMENT_RE = re.compile(r"[^A-Za-z0-9._+-]+")
class LockContentionError(RuntimeError):
"""Raised when an exclusive per-issue lock cannot be acquired."""
def default_lock_dir() -> str: def default_lock_dir() -> str:
override = (os.environ.get("GITEA_ISSUE_LOCK_DIR") or "").strip() raw = (os.environ.get(LOCK_DIR_ENV) or DEFAULT_LOCK_DIR).strip()
if override: return raw or DEFAULT_LOCK_DIR
return override
cache_dir = (os.environ.get("GITEA_MCP_SESSION_STATE_DIR") or "").strip()
if cache_dir: def _sanitize_segment(value: str) -> str:
return os.path.join(cache_dir, "locks") text = (value or "").strip()
user_cache = os.path.expanduser("~/.cache") if not text:
return os.path.join(user_cache, "gitea-tools", "locks") return "_"
return _SAFE_SEGMENT_RE.sub("_", text)
def lock_key(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
) -> str:
return "-".join(
_sanitize_segment(part)
for part in (remote, org, repo, str(issue_number))
)
def lock_file_path(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
lock_dir: str | None = None,
) -> str:
root = (lock_dir or default_lock_dir()).strip()
return os.path.join(root, f"{lock_key(remote=remote, org=org, repo=repo, issue_number=issue_number)}.json")
def session_pointer_path(lock_dir: str | None = None) -> str: def session_pointer_path(lock_dir: str | None = None) -> str:
@@ -42,60 +83,78 @@ def flock_path(json_path: str) -> str:
return f"{json_path}.lock" return f"{json_path}.lock"
@contextlib.contextmanager def is_process_alive(pid: int | None) -> bool:
if not pid or pid <= 0:
return False
try:
os.kill(int(pid), 0)
return True
except OSError as exc:
return exc.errno != errno.ESRCH
except (TypeError, ValueError):
return False
@contextmanager
def _exclusive_file_lock(lock_path: str): def _exclusive_file_lock(lock_path: str):
import fcntl os.makedirs(os.path.dirname(lock_path) or ".", exist_ok=True)
fd = os.open(lock_path, os.O_CREAT | os.O_RDWR, 0o600)
_ensure_lock_dir(os.path.dirname(lock_path)) try:
with open(lock_path, "a+") as fh:
try: try:
fcntl.flock(fh.fileno(), fcntl.LOCK_EX) fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
yield except BlockingIOError as exc:
raise LockContentionError(
f"could not acquire exclusive lock on '{lock_path}'"
) from exc
yield fd
finally:
try:
fcntl.flock(fd, fcntl.LOCK_UN)
finally: finally:
try: os.close(fd)
fcntl.flock(fh.fileno(), fcntl.LOCK_UN)
except Exception:
pass
def lock_file_path(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
lock_dir: str | None = None,
) -> str:
root = (lock_dir or default_lock_dir()).strip()
slug = f"{remote}-{org}-{repo}-issue-{issue_number}.json"
safe_slug = re.sub(r"[^A-Za-z0-9._-]", "_", slug)
return os.path.join(root, safe_slug)
def read_lock_file(path: str) -> dict[str, Any] | None: def read_lock_file(path: str) -> dict[str, Any] | None:
if not path or not os.path.isfile(path): lock_path = (path or "").strip()
if not lock_path or not os.path.exists(lock_path):
return None return None
try: try:
with open(path, "r", encoding="utf-8") as fh: with open(lock_path, encoding="utf-8") as handle:
data = json.load(fh) data = json.load(handle)
if isinstance(data, dict): except (OSError, json.JSONDecodeError):
return data return None
except Exception: return data if isinstance(data, dict) else None
pass
return None
def save_lock_file(path: str, data: dict[str, Any]) -> None: def save_lock_file(path: str, data: dict[str, Any]) -> None:
_ensure_lock_dir(os.path.dirname(path)) lock_path = (path or "").strip()
tmp_path = f"{path}.tmp.{os.getpid()}" if not lock_path:
with open(tmp_path, "w", encoding="utf-8") as fh: raise ValueError("lock path is required (fail closed)")
json.dump(data, fh, indent=2) parent = os.path.dirname(lock_path) or "."
fh.flush() os.makedirs(parent, mode=0o700, exist_ok=True)
os.fsync(fh.fileno()) payload = json.dumps(data, indent=2, sort_keys=True) + "\n"
os.replace(tmp_path, path) fd, temp_path = tempfile.mkstemp(prefix=".lock-", suffix=".json", dir=parent)
try:
with os.fdopen(fd, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
os.replace(temp_path, lock_path)
finally:
if os.path.exists(temp_path):
try:
os.remove(temp_path)
except OSError:
pass
def lock_generation(lock: dict[str, Any] | None) -> int: def lock_generation(lock: dict[str, Any] | None) -> int:
"""Monotonic write counter for a durable lock record (#772 AC5).
Absent or unusable values read as ``0`` so a lock written before generations
existed still participates in compare-and-swap: its first recovery expects
``0`` and writes ``1``.
"""
if not isinstance(lock, dict): if not isinstance(lock, dict):
return 0 return 0
try: try:
@@ -112,6 +171,16 @@ def bind_session_lock(
renewal_sanctioned: bool = False, renewal_sanctioned: bool = False,
recovery_sanctioned: bool = False, recovery_sanctioned: bool = False,
) -> str: ) -> str:
"""Persist a keyed lock and bind it to the current process session.
``expected_generation`` turns the write into a compare-and-swap (#772 AC5).
Recovery decides it may take over a claim by reading the durable lock, but
that read and this write are separate steps; without a CAS two replacement
sessions can both observe the same dead owner, both pass assessment, and
both write — the second silently clobbering the first. Passing the
generation observed at assessment time makes exactly one of them win: the
loser's expectation no longer matches and it fails closed.
"""
remote = str(lock_data.get("remote") or "") remote = str(lock_data.get("remote") or "")
org = str(lock_data.get("org") or "") org = str(lock_data.get("org") or "")
repo = str(lock_data.get("repo") or "") repo = str(lock_data.get("repo") or "")
@@ -160,33 +229,31 @@ def bind_session_lock(
) )
if lease_block: if lease_block:
raise RuntimeError(lease_block) raise RuntimeError(lease_block)
# #772 AC5: compare-and-swap inside the same critical section that
current_gen = lock_generation(existing) # already serializes writers, so the check and the write cannot be
if expected_generation is not None and current_gen != expected_generation: # separated by another session's successful recovery.
current_generation = lock_generation(existing)
if (
expected_generation is not None
and current_generation != expected_generation
):
raise RuntimeError( raise RuntimeError(
f"compare-and-swap generation mismatch on issue #{issue_number}: " f"Issue #{issue_number} lock generation changed: expected "
f"expected {expected_generation}, observed {current_gen} (fail closed)" f"{expected_generation}, found {current_generation}; another "
"session already recovered or replaced this claim (fail closed)"
) )
record["lock_generation"] = current_generation + 1
record["lock_generation"] = current_gen + 1
if not record.get("lock_provenance"):
provenance_source = (
issue_lock_provenance.SOURCE_RECOVER_DIRTY_ORPHANED
if recovery_sanctioned
else issue_lock_provenance.SOURCE_LOCK_ISSUE
)
record["lock_provenance"] = issue_lock_provenance.build_lock_provenance(
issue_number=issue_number,
source=provenance_source,
branch_name=record.get("branch_name"),
worktree_path=record.get("worktree_path"),
generation=record["lock_generation"],
)
save_lock_file(path, record) save_lock_file(path, record)
save_lock_file(session_pointer_path(root), pointer) save_lock_file(session_pointer_path(root), pointer)
except Exception: except LockContentionError as exc:
raise competing = read_lock_file(path)
if competing:
owner_pid = competing.get("session_pid") or competing.get("pid")
raise RuntimeError(
f"Issue #{issue_number} lock contention: {exc}; competing owner "
f"pid={owner_pid} (fail closed)"
) from exc
raise RuntimeError(f"Issue #{issue_number} lock contention: {exc} (fail closed)") from exc
return path return path
@@ -224,18 +291,21 @@ def iter_lock_files(lock_dir: str | None = None) -> list[str]:
root = (lock_dir or default_lock_dir()).strip() root = (lock_dir or default_lock_dir()).strip()
if not os.path.isdir(root): if not os.path.isdir(root):
return [] return []
out: list[str] = [] paths: list[str] = []
for entry in os.listdir(root): for name in os.listdir(root):
if entry.endswith(".json") and not entry.startswith("session-"): if not name.endswith(".json") or name.startswith("session-"):
out.append(os.path.join(root, entry)) continue
return sorted(out) paths.append(os.path.join(root, name))
return sorted(paths)
def find_live_lock_for_branch( def find_lock_for_branch(
*,
remote: str,
org: str,
repo: str,
branch_name: str, branch_name: str,
lock_dir: str | None = None, lock_dir: str | None = None,
*,
now: datetime | None = None,
) -> dict[str, Any] | None: ) -> dict[str, Any] | None:
target = (branch_name or "").strip() target = (branch_name or "").strip()
if not target: if not target:
@@ -244,42 +314,29 @@ def find_live_lock_for_branch(
lock = read_lock_file(path) lock = read_lock_file(path)
if not lock: if not lock:
continue continue
lock_branch = str(lock.get("branch_name") or "").strip() if (
if lock_branch == target and is_lease_live(lock, now=now): str(lock.get("remote") or "") == remote
and str(lock.get("org") or "") == org
and str(lock.get("repo") or "") == repo
and str(lock.get("branch_name") or "").strip() == target
):
lock = dict(lock)
lock.setdefault("lock_file_path", path)
return lock return lock
return None return None
def is_process_alive(pid: int | None) -> bool:
if pid is None or pid <= 0:
return False
try:
os.kill(pid, 0)
return True
except OSError:
return False
def _lease_now(now: datetime | None = None) -> datetime: def _lease_now(now: datetime | None = None) -> datetime:
if now is not None: return now or datetime.now(timezone.utc)
if now.tzinfo is None:
return now.replace(tzinfo=timezone.utc)
return now.astimezone(timezone.utc)
return datetime.now(timezone.utc)
def _parse_lease_timestamp(text: str | None) -> datetime | None: def _parse_lease_timestamp(value: str | None) -> datetime | None:
text = (value or "").strip()
if not text: if not text:
return None return None
raw = str(text).strip()
if not raw:
return None
try: try:
dt = datetime.fromisoformat(raw) return datetime.fromisoformat(text.replace("Z", "+00:00")).astimezone(timezone.utc)
if dt.tzinfo is None: except ValueError:
return dt.replace(tzinfo=timezone.utc)
return dt.astimezone(timezone.utc)
except Exception:
return None return None
@@ -308,6 +365,7 @@ def assess_lock_freshness(
*, *,
now: datetime | None = None, now: datetime | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Classify a lock as live, expired, stale, or absent."""
current = _lease_now(now) current = _lease_now(now)
if not lock_data: if not lock_data:
return { return {
@@ -349,6 +407,10 @@ def assess_lock_freshness(
"pid_missing": pid_missing, "pid_missing": pid_missing,
} }
# #860: a PID-less lock must never be considered live merely because
# expiration / heartbeat fields are absent. Missing PID is insufficient
# evidence of a live owner; treat as malformed/stale so recovery routes
# can evaluate corroborating pins instead of blocking on a false live flag.
if pid_missing: if pid_missing:
return { return {
"status": "malformed", "status": "malformed",
@@ -369,31 +431,87 @@ def assess_lock_freshness(
"status": "stale", "status": "stale",
"live": False, "live": False,
"stale": True, "stale": True,
"reason": f"owner pid {pid_int} is dead", "reason": f"owner pid {pid_int} is not alive",
"pid_alive": False, "pid_alive": False,
"pid_missing": False, "pid_missing": False,
"pid": pid_int,
}
raw_status = (lock_data.get("status") or "").strip().lower()
if raw_status in {"active", "acquired", "locked"}:
return {
"status": "active",
"live": True,
"stale": False,
"reason": f"lock active (pid {pid_int})",
"pid_alive": True,
"pid_missing": False,
"pid": pid_int,
} }
return { return {
"status": raw_status or "unknown", "status": "live",
"live": False, "live": True,
"stale": True, "stale": False,
"reason": f"unrecognized lock status '{raw_status}'", "reason": "lock heartbeat and lease are fresh",
"pid_alive": pid_alive, "pid_alive": pid_alive,
"pid_missing": False, "pid_missing": False,
"heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None,
"expires_at": expires_at.isoformat() if expires_at else None,
}
def _same_realpath(left: str | None, right: str | None) -> bool:
if not left or not right:
return False
try:
return os.path.realpath(left) == os.path.realpath(right)
except OSError:
return left == right
def assess_expired_lock_reclaim(
existing_lock: dict[str, Any] | None,
*,
now: datetime | None = None,
) -> dict[str, Any]:
"""Decide whether an expired/stale author issue lock may be reclaimed (#601).
Required proof (any reclaim of non-live lock):
* lease not live (expired or dead pid)
* owner process dead OR worktree missing
* no force-delete of live foreign ownership
"""
if not existing_lock:
return {
"reclaim_allowed": True,
"reasons": ["no existing lock"],
"freshness": assess_lock_freshness(None, now=now),
}
freshness = assess_lock_freshness(existing_lock, now=now)
if freshness.get("live"):
return {
"reclaim_allowed": False,
"reasons": ["lock is still live; cannot reclaim (fail closed)"],
"freshness": freshness,
}
pid = existing_lock.get("session_pid")
if pid is None:
pid = existing_lock.get("pid")
dead = not is_process_alive(pid) if pid is not None else True
wt = str(existing_lock.get("worktree_path") or "")
missing_wt = (not wt) or (not os.path.isdir(os.path.realpath(wt)))
if not (dead or missing_wt):
return {
"reclaim_allowed": False,
"reasons": [
"expired/stale lock still has live owner pid and present worktree; "
"recovery review required (fail closed)"
],
"freshness": freshness,
"owner_pid_dead": dead,
"worktree_missing": missing_wt,
}
return {
"reclaim_allowed": True,
"reasons": [
"non-live lock with dead process and/or missing worktree; "
"sanctioned reclaim allowed"
],
"freshness": freshness,
"owner_pid_dead": dead,
"worktree_missing": missing_wt,
"prior_branch": existing_lock.get("branch_name"),
"prior_worktree": existing_lock.get("worktree_path"),
"prior_pid": pid,
} }
@@ -403,98 +521,299 @@ def assess_same_issue_lease_conflict(
issue_number: int, issue_number: int,
branch_name: str, branch_name: str,
worktree_path: str, worktree_path: str,
now: datetime | None = None, operation_type: str = AUTHOR_ISSUE_WORK_LEASE,
renewal_sanctioned: bool = False, renewal_sanctioned: bool = False,
recovery_sanctioned: bool = False, recovery_sanctioned: bool = False,
now: datetime | None = None,
) -> str | None: ) -> str | None:
"""Return a fail-closed error when a competing live lease blocks acquisition.
``renewal_sanctioned`` is set only when
``issue_lock_renewal.assess_exact_owner_lease_renewal`` has already proven,
from the durable lock plus live server-side observation, that this session
is the exact recorded owner of an *expired* lease (#760). It is never a
caller-supplied parameter of any MCP tool (#760 AC14): the server computes
it and passes it down. Left False, every pre-existing disposition is
unchanged.
"""
if not existing_lock: if not existing_lock:
return None return None
freshness = assess_lock_freshness(existing_lock, now=now)
if not freshness["live"]: existing_issue = existing_lock.get("issue_number")
return None lease = existing_lock.get("work_lease")
if renewal_sanctioned or recovery_sanctioned: existing_operation = (
lease.get("operation_type")
if isinstance(lease, dict)
else AUTHOR_ISSUE_WORK_LEASE
)
if existing_issue != issue_number or existing_operation != operation_type:
return None return None
existing_wt = str(existing_lock.get("worktree_path") or "").strip() existing_branch = existing_lock.get("branch_name")
target_wt = (worktree_path or "").strip() existing_worktree = existing_lock.get("worktree_path")
if existing_wt and target_wt: same_owner = (
try: existing_branch == branch_name
same_wt = os.path.realpath(existing_wt) == os.path.realpath(target_wt) and _same_realpath(str(existing_worktree or ""), worktree_path)
except Exception: )
same_wt = existing_wt == target_wt if recovery_sanctioned and existing_issue == issue_number and existing_branch == branch_name:
if not same_wt: return None
owner_pid = existing_lock.get("session_pid") or existing_lock.get("pid") if is_lease_expired(existing_lock, now=now):
return ( # #760 AC1/AC2: exact-owner renewal is a different disposition from
f"Issue #{issue_number} already has an active author_issue_work lease " # foreign takeover and is evaluated first. Before this, both branches
f"from worktree '{existing_wt}' (pid={owner_pid}; fail closed)" # below returned unconditionally, so the same_owner allowance further
) # down was unreachable for every expired lease — an owner could never
# renew its own lock once the wall clock passed, no matter how complete
existing_br = str(existing_lock.get("branch_name") or "").strip() # its ownership evidence. Requires BOTH the locally recomputed
target_br = (branch_name or "").strip() # same_owner match and the server-proven renewal waiver; either alone is
if existing_br and target_br and existing_br != target_br: # insufficient.
if same_owner and renewal_sanctioned:
return None
reclaim = assess_expired_lock_reclaim(existing_lock, now=now)
if reclaim.get("reclaim_allowed"):
# #601: expired + dead pid / missing worktree may be reclaimed
# through the normal lock path (sanctioned overwrite).
return None
return ( return (
f"Issue #{issue_number} already has an active lease on branch '{existing_br}' " f"Issue #{issue_number} has an expired {operation_type} lease on "
f"(cannot lock for branch '{target_br}'; fail closed)" f"branch '{existing_branch}' from worktree '{existing_worktree}'. "
"Recovery review is required before takeover (fail closed)"
) )
return None if same_owner:
return None
return (
f"Issue #{issue_number} already has an active {operation_type} lease on "
f"branch '{existing_branch}' from worktree '{existing_worktree}' "
"(fail closed)"
)
def _lock_claimant(lock: dict[str, Any] | None) -> dict[str, str]:
if not isinstance(lock, dict):
return {}
claimant = lock.get("claimant")
if not isinstance(claimant, dict):
lease = lock.get("work_lease")
claimant = lease.get("claimant") if isinstance(lease, dict) else None
if not isinstance(claimant, dict):
return {}
return {
"username": str(claimant.get("username") or ""),
"profile": str(claimant.get("profile") or ""),
}
def assess_foreign_lock_overwrite( def assess_foreign_lock_overwrite(
existing_lock: dict[str, Any] | None, existing_lock: dict[str, Any] | None,
proposed_lock: dict[str, Any], incoming_lock: dict[str, Any],
*, *,
now: datetime | None = None,
recovery_sanctioned: bool = False, recovery_sanctioned: bool = False,
now: datetime | None = None,
) -> str | None: ) -> str | None:
"""Block writes that would clobber an unrelated live lease on the same key."""
if not existing_lock: if not existing_lock:
return None return None
freshness = assess_lock_freshness(existing_lock, now=now)
ex_claimant = ( same_issue = existing_lock.get("issue_number") == incoming_lock.get("issue_number")
existing_lock.get("claimant") same_branch = existing_lock.get("branch_name") == incoming_lock.get("branch_name")
or existing_lock.get("user") same_worktree = _same_realpath(
or existing_lock.get("username") str(existing_lock.get("worktree_path") or ""),
or existing_lock.get("owner") str(incoming_lock.get("worktree_path") or ""),
or ""
) )
prop_claimant = ( if same_issue and same_branch and same_worktree:
proposed_lock.get("claimant") return None
or proposed_lock.get("user")
or proposed_lock.get("username") existing_claimant = _lock_claimant(existing_lock)
or proposed_lock.get("owner") incoming_claimant = _lock_claimant(incoming_lock)
or "" same_claimant = (
) bool(existing_claimant.get("username"))
ex_profile = ( and existing_claimant.get("username") == incoming_claimant.get("username")
existing_lock.get("profile") and existing_claimant.get("profile") == incoming_claimant.get("profile")
or existing_lock.get("profile_name")
or existing_lock.get("execution_profile")
or ""
)
prop_profile = (
proposed_lock.get("profile")
or proposed_lock.get("profile_name")
or proposed_lock.get("execution_profile")
or ""
) )
same_claimant = bool( if recovery_sanctioned and same_issue and same_branch and same_claimant:
ex_claimant and prop_claimant and ex_claimant == prop_claimant return None
)
same_profile = bool( if not is_lease_live(existing_lock, now=now):
ex_profile and prop_profile and ex_profile == prop_profile # #860 F8: A non-live or PID-less lock still blocks foreign overwrite
# unless same claimant or sanctioned reclaim is proven.
if not same_claimant and same_issue:
reclaim = assess_expired_lock_reclaim(existing_lock, now=now)
if not reclaim.get("reclaim_allowed"):
return (
"Refusing foreign overwrite of non-live issue lock "
f"(issue #{existing_lock.get('issue_number')}, owner '{existing_claimant.get('username')}') "
"without sanctioned reclaim proof (fail closed)"
)
return None
return (
"Refusing to overwrite a live foreign issue lock "
f"(issue #{existing_lock.get('issue_number')}, "
f"branch '{existing_lock.get('branch_name')}', "
f"worktree '{existing_lock.get('worktree_path')}') (fail closed)"
) )
if not same_claimant and not recovery_sanctioned:
return (
f"Foreign lock overwrite refused: existing lock belongs to claimant "
f"'{ex_claimant or 'unknown'}' (proposed: '{prop_claimant or 'unknown'}'); "
f"foreign locks may only be overwritten through explicit sanctioned "
f"recovery (fail closed)"
)
if freshness["live"] and not (same_claimant or same_profile) and not recovery_sanctioned: def find_live_lock_for_branch(
return ( branch_name: str,
f"Foreign lock overwrite refused: existing lock is live " lock_dir: str | None = None,
f"(status={freshness['status']}) and belongs to '{ex_claimant or 'unknown'}'" ) -> dict[str, Any] | None:
) target = (branch_name or "").strip()
if not target:
return None
for path in iter_lock_files(lock_dir):
lock = read_lock_file(path)
if not lock:
continue
if str(lock.get("branch_name") or "").strip() != target:
continue
if not is_lease_live(lock):
continue
record = dict(lock)
record.setdefault("lock_file_path", path)
return record
return None return None
def resolve_locked_branch_for_session(
branch_name: str | None = None,
lock_dir: str | None = None,
) -> str:
if branch_name:
lock = find_live_lock_for_branch(branch_name, lock_dir)
if lock:
return str(lock.get("branch_name") or "")
lock = read_session_issue_lock(lock_dir)
return str((lock or {}).get("branch_name") or "")
def has_active_issue_lock(
branch: str,
*,
lock_dir: str | None = None,
) -> bool:
target = (branch or "").strip()
if not target:
return False
for path in iter_lock_files(lock_dir):
lock = read_lock_file(path)
if not lock:
continue
if str(lock.get("branch_name") or "").strip() != target:
continue
if is_lease_live(lock):
return True
return False
def verify_lock_for_mutation(
lock_data: dict[str, Any] | None,
*,
issue_number: int | None = None,
branch_name: str | None = None,
worktree_path: str | None = None,
) -> dict[str, Any]:
"""Re-check lock ownership immediately before a mutation (#438)."""
reasons: list[str] = []
if not lock_data:
return {"proven": False, "block": True, "reasons": ["issue lock is missing (fail closed)"]}
freshness = assess_lock_freshness(lock_data)
if not freshness["live"]:
reasons.append(f"issue lock is not live: {freshness['reason']} (fail closed)")
if issue_number is not None and lock_data.get("issue_number") != issue_number:
reasons.append(
f"issue lock targets #{lock_data.get('issue_number')}, expected #{issue_number} (fail closed)"
)
if branch_name is not None and lock_data.get("branch_name") != branch_name:
reasons.append(
f"issue lock branch '{lock_data.get('branch_name')}' does not match "
f"'{branch_name}' (fail closed)"
)
if worktree_path is not None:
locked = os.path.realpath(str(lock_data.get("worktree_path") or ""))
declared = os.path.realpath(worktree_path)
if locked != declared:
reasons.append(
f"issue lock worktree '{locked}' does not match declared '{declared}' (fail closed)"
)
return {
"proven": not reasons,
"block": bool(reasons),
"reasons": reasons,
"freshness": freshness,
"lock_proof": format_lock_proof(lock_data, freshness=freshness),
}
def list_live_locks(
*,
lock_dir: str | None = None,
now: datetime | None = None,
) -> list[dict[str, Any]]:
"""Return live per-issue locks for queue visibility."""
live: list[dict[str, Any]] = []
for path in iter_lock_files(lock_dir):
record = read_lock_file(path)
if not record:
continue
freshness = assess_lock_freshness(record, now=now)
if not freshness["live"]:
continue
live.append(
{
"issue_number": record.get("issue_number"),
"branch_name": record.get("branch_name"),
"remote": record.get("remote"),
"org": record.get("org"),
"repo": record.get("repo"),
"worktree_path": record.get("worktree_path"),
"pid": record.get("session_pid") or record.get("pid"),
"claimant": (
record.get("claimant")
or (record.get("work_lease") or {}).get("claimant")
),
"freshness": freshness,
"lock_path": record.get("lock_file_path") or path,
}
)
return live
def format_lock_proof(
lock_data: dict[str, Any] | None,
*,
freshness: dict[str, Any] | None = None,
competing_live_locks: list[dict[str, Any]] | None = None,
released: bool | None = None,
) -> str:
"""Canonical issue-lock proof string for final reports."""
if not lock_data:
return "issue lock proof: not acquired"
fresh = freshness or assess_lock_freshness(lock_data)
owner = lock_data.get("claimant") or {}
if not owner and isinstance(lock_data.get("work_lease"), dict):
owner = lock_data["work_lease"].get("claimant") or {}
parts = [
"issue lock proof:",
f"acquired issue #{lock_data.get('issue_number')}",
f"branch {lock_data.get('branch_name')}",
f"owner {owner.get('profile') or 'unknown'}",
f"pid {lock_data.get('session_pid') or lock_data.get('pid')}",
f"freshness {fresh.get('status')}",
]
if competing_live_locks is not None:
parts.append(
"no competing live lock"
if not competing_live_locks
else f"competing live locks {len(competing_live_locks)}"
)
if released is True:
parts.append("lock released")
elif released is False:
parts.append("lock retained")
return "; ".join(parts)
+11
View File
@@ -41,6 +41,17 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.issue.comment", "permission": "gitea.issue.comment",
"role": "author", "role": "author",
}, },
# #864: dirty-preserving same-claimant author-session rebind (dead owner PID).
# Author MCP tool path. Reconciler execute is gated inside the tool via
# authorize_reconciler_execute + role_kind checks (not this map entry).
"rebind_dirty_same_claimant_author_session": {
"permission": "gitea.issue.comment",
"role": "author",
},
"gitea_rebind_dirty_same_claimant_author_session": {
"permission": "gitea.issue.comment",
"role": "author",
},
"set_issue_labels": { "set_issue_labels": {
"permission": "gitea.issue.comment", "permission": "gitea.issue.comment",
"role": "author", "role": "author",
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,551 @@
"""Merged-PR awareness for the worktree cleanup audit (#858).
Before #858 an ``issue_work`` worktree could never leave ``active_issue_work``:
the audit had no PR linkage at all (``pr_number`` was structurally ``None``)
and its only route to ``clean_stale_removable`` was a TTL derived from a
``last_used_at`` that nothing ever populated. A merged, clean, unprotected
worktree was therefore reported as active work forever, disagreeing with the
PR-scoped reconciler.
These tests use fabricated temporary repositories and synthetic PR records
only. Nothing here removes a worktree or deletes a branch.
"""
import os
import subprocess
import sys
import tempfile
import unittest
from unittest.mock import patch
sys.path.insert(0, str(__import__("pathlib").Path(__file__).resolve().parent.parent))
import merged_cleanup_reconcile as mcr # noqa: E402
import worktree_cleanup_audit as wca # noqa: E402
MERGED_BRANCH = "feat/issue-777-timeline"
MERGED_PATH = "/repo/branches/issue-777-timeline"
HEAD_SHA = "a" * 40
def _pr(number, branch, *, merged=True, sha=HEAD_SHA, state=None):
"""Synthetic Gitea PR payload."""
return {
"number": number,
"head": {"ref": branch, "sha": sha},
"merged_at": "2026-07-24T01:00:00Z" if merged else None,
"state": state or ("closed" if merged else "open"),
}
def _porcelain(*entries):
out = []
for path, branch, sha in entries:
out.append(f"worktree {path}")
out.append(f"HEAD {sha}")
if branch is None:
out.append("detached")
else:
out.append(f"branch refs/heads/{branch}")
out.append("")
return "\n".join(out)
class _AuditHarness(unittest.TestCase):
"""Runs audit_branches_directory over a fabricated worktree listing."""
PORCELAIN = _porcelain(
("/repo", "master", "f" * 40),
(MERGED_PATH, MERGED_BRANCH, HEAD_SHA),
)
def run_audit(self, *, dirty_paths=(), contained=True, **kwargs):
def fake_dirty(path):
if path in dirty_paths:
return {"exists": True, "dirty": True, "dirty_files": [" M x.py"]}
return {"exists": True, "dirty": False, "dirty_files": []}
with patch.object(
wca, "list_worktrees",
return_value=wca.parse_worktree_porcelain(self.PORCELAIN),
), patch.object(
wca, "read_worktree_dirty", side_effect=fake_dirty
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=contained
):
report = wca.audit_branches_directory("/repo", **kwargs)
return {wt["path"]: wt for wt in report["worktrees"]}, report
def merged_audit(self, **kwargs):
kwargs.setdefault("pr_index", wca.build_pr_index([_pr(849, MERGED_BRANCH)]))
kwargs.setdefault("master_ref", "prgs/master")
return self.run_audit(**kwargs)
class TestMergedWorktreeBecomesRemovable(_AuditHarness):
def test_clean_merged_issue_worktree_is_linked_and_removable(self):
by_path, report = self.merged_audit()
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_CLEAN_STALE_REMOVABLE)
self.assertTrue(entry["removable"])
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_MERGED)
self.assertEqual(entry["merged_pr_cleanup"]["block_reasons"], [])
self.assertIn(MERGED_PATH, [c["path"] for c in report["removable_candidates"]])
def test_pr_number_populated_from_authoritative_linkage(self):
by_path, _ = self.merged_audit()
self.assertEqual(by_path[MERGED_PATH]["pr_number"], 849)
def test_regression_without_pr_evidence_stays_active_issue_work(self):
"""The pre-#858 behaviour, still correct when no PR state is supplied."""
by_path, _ = self.run_audit()
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
self.assertIsNone(entry["pr_number"])
class TestProtectiveSignalsSurvive(_AuditHarness):
def test_open_pr_worktree_is_not_removable(self):
index = wca.build_pr_index([_pr(900, MERGED_BRANCH, merged=False)])
by_path, _ = self.run_audit(
pr_index=index,
master_ref="prgs/master",
open_pr_branches={MERGED_BRANCH},
)
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_OPEN_PR)
self.assertFalse(entry["removable"])
# linkage still reports the owning PR, it just is not merge proof
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_OPEN)
self.assertEqual(entry["pr_number"], 900)
def test_dirty_tracked_worktree_is_not_removable(self):
by_path, _ = self.merged_audit(dirty_paths=(MERGED_PATH,))
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_DIRTY_LOCAL)
self.assertFalse(entry["removable"])
self.assertIn(
"worktree has uncommitted changes",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_untracked_only_worktree_is_not_removable(self):
"""``git status --porcelain`` reports untracked files as dirty too."""
def untracked(path):
if path == MERGED_PATH:
return {"exists": True, "dirty": True, "dirty_files": ["?? scratch.txt"]}
return {"exists": True, "dirty": False, "dirty_files": []}
with patch.object(
wca, "list_worktrees",
return_value=wca.parse_worktree_porcelain(self.PORCELAIN),
), patch.object(
wca, "read_worktree_dirty", side_effect=untracked
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory(
"/repo",
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
)
entry = {wt["path"]: wt for wt in report["worktrees"]}[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_DIRTY_LOCAL)
self.assertFalse(entry["removable"])
def test_active_lease_by_issue_number_is_protective(self):
by_path, _ = self.merged_audit(leased_issue_numbers={777})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_active_lease"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_active_lease_by_branch_is_protective(self):
by_path, _ = self.merged_audit(leased_branches={MERGED_BRANCH})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_active_lease"])
self.assertFalse(entry["removable"])
def test_active_issue_lock_is_protective(self):
by_path, _ = self.merged_audit(active_issue_branches={MERGED_BRANCH})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_active_issue_lock"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_live_session_worktree_is_protective(self):
by_path, _ = self.merged_audit(live_session_paths={MERGED_PATH})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_live_session"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_head_not_contained_in_master_is_not_removable(self):
by_path, _ = self.merged_audit(contained=False)
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
self.assertIn(
"worktree head is not contained in authoritative master "
"(unmerged commits remain)",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_unknown_containment_fails_closed(self):
by_path, _ = self.merged_audit(contained=None)
entry = by_path[MERGED_PATH]
self.assertFalse(entry["removable"])
self.assertIn(
"containment of the worktree head in master is unknown",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_missing_master_ref_fails_closed(self):
by_path, _ = self.run_audit(
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)])
)
self.assertFalse(by_path[MERGED_PATH]["removable"])
def test_unmerged_owning_pr_is_not_removable(self):
index = wca.build_pr_index([_pr(901, MERGED_BRANCH, merged=False)])
by_path, _ = self.run_audit(pr_index=index, master_ref="prgs/master")
entry = by_path[MERGED_PATH]
self.assertFalse(entry["removable"])
self.assertIn(
"owning PR #901 is not merged",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_control_checkout_is_never_removable(self):
by_path, _ = self.merged_audit()
control = by_path["/repo"]
self.assertTrue(control["is_protected"])
self.assertEqual(control["classification"], wca.CLASS_UNSAFE_UNKNOWN)
self.assertFalse(control["removable"])
def test_control_checkout_not_removable_even_if_linked_and_merged(self):
"""A merged PR on the control checkout must not unlock removal."""
porcelain = _porcelain(("/repo", MERGED_BRANCH, HEAD_SHA))
with patch.object(
wca, "list_worktrees", return_value=wca.parse_worktree_porcelain(porcelain)
), patch.object(
wca, "read_worktree_dirty",
return_value={"exists": True, "dirty": False, "dirty_files": []},
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory(
"/repo",
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
)
entry = report["worktrees"][0]
self.assertEqual(entry["classification"], wca.CLASS_UNSAFE_UNKNOWN)
self.assertFalse(entry["removable"])
class TestAmbiguousLinkageFailsClosed(_AuditHarness):
def test_competing_prs_on_one_branch_fail_closed(self):
index = wca.build_pr_index(
[_pr(849, MERGED_BRANCH), _pr(860, MERGED_BRANCH)]
)
by_path, _ = self.run_audit(pr_index=index, master_ref="prgs/master")
entry = by_path[MERGED_PATH]
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_AMBIGUOUS)
self.assertIsNone(entry["pr_number"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_merged_plus_open_pr_on_one_branch_fails_closed(self):
index = wca.build_pr_index(
[_pr(849, MERGED_BRANCH), _pr(861, MERGED_BRANCH, merged=False)]
)
by_path, _ = self.run_audit(pr_index=index, master_ref="prgs/master")
entry = by_path[MERGED_PATH]
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_AMBIGUOUS)
self.assertFalse(entry["removable"])
def test_no_owning_pr_fails_closed(self):
by_path, _ = self.run_audit(
pr_index=wca.build_pr_index([_pr(849, "feat/other-branch")]),
master_ref="prgs/master",
)
entry = by_path[MERGED_PATH]
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_NONE)
self.assertFalse(entry["removable"])
def test_malformed_pr_records_are_dropped_not_guessed(self):
index = wca.build_pr_index(
[
{"number": None, "head": {"ref": MERGED_BRANCH}},
{"number": 5, "head": {}},
{"number": "not-an-int", "head": {"ref": MERGED_BRANCH}},
]
)
self.assertEqual(index, {})
self.assertEqual(
wca.resolve_owning_pr(branch=MERGED_BRANCH, pr_index=index)["status"],
wca.LINKAGE_NONE,
)
def test_detached_worktree_has_no_branch_linkage(self):
self.assertEqual(
wca.resolve_owning_pr(branch=None, pr_index={})["status"],
wca.LINKAGE_UNKNOWN,
)
class TestUnrelatedClassificationsUnchanged(unittest.TestCase):
"""Non-issue_work worktrees keep their pre-#858 classifications."""
PORCELAIN = _porcelain(
("/repo", "master", "f" * 40),
("/repo/branches/review-pr42", "review-pr42", "2" * 40),
("/repo/branches/baseline-master-x", "baseline-master-x", "3" * 40),
("/repo/branches/conflict-fix-pr50", "conflict-fix-pr50", "4" * 40),
("/repo/branches/review-pr99", None, "5" * 40),
)
def _audit(self, **kwargs):
with patch.object(
wca, "list_worktrees",
return_value=wca.parse_worktree_porcelain(self.PORCELAIN),
), patch.object(
wca, "read_worktree_dirty",
return_value={"exists": True, "dirty": False, "dirty_files": []},
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory("/repo", **kwargs)
return {wt["path"]: wt for wt in report["worktrees"]}
def test_classifications_identical_with_and_without_pr_evidence(self):
without = self._audit()
with_evidence = self._audit(
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
)
self.assertEqual(
{p: e["classification"] for p, e in without.items()},
{p: e["classification"] for p, e in with_evidence.items()},
)
def test_lease_on_issue_does_not_capture_similarly_named_scratch_trees(self):
"""A lease on issue 777 protects issue work, not baseline/review trees."""
porcelain = _porcelain(
("/repo/branches/baseline-master-issue-777", "baseline-issue-777", "7" * 40),
("/repo/branches/issue-777-timeline", MERGED_BRANCH, HEAD_SHA),
)
with patch.object(
wca, "list_worktrees", return_value=wca.parse_worktree_porcelain(porcelain)
), patch.object(
wca, "read_worktree_dirty",
return_value={"exists": True, "dirty": False, "dirty_files": []},
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory(
"/repo",
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
leased_issue_numbers={777},
)
by_path = {wt["path"]: wt for wt in report["worktrees"]}
baseline = by_path["/repo/branches/baseline-master-issue-777"]
self.assertFalse(baseline["has_active_lease"])
self.assertEqual(baseline["classification"], wca.CLASS_CLEAN_STALE_REMOVABLE)
issue_work = by_path["/repo/branches/issue-777-timeline"]
self.assertTrue(issue_work["has_active_lease"])
self.assertFalse(issue_work["removable"])
def test_review_and_baseline_still_removable(self):
by_path = self._audit(
pr_index=wca.build_pr_index([]), master_ref="prgs/master"
)
self.assertEqual(
by_path["/repo/branches/review-pr42"]["classification"],
wca.CLASS_CLEAN_STALE_REMOVABLE,
)
self.assertEqual(
by_path["/repo/branches/baseline-master-x"]["classification"],
wca.CLASS_CLEAN_STALE_REMOVABLE,
)
self.assertEqual(
by_path["/repo/branches/review-pr99"]["classification"],
wca.CLASS_DETACHED_REVIEW_LEFTOVER,
)
def test_conflict_fix_ttl_behaviour_unchanged(self):
"""conflict_fix still needs only TTL expiry; #858 did not touch it."""
self.assertEqual(
wca.classify_worktree(
workflow_type=wca.WORKFLOW_CONFLICT_FIX,
is_dirty=False,
ttl_expired=True,
),
wca.CLASS_CLEAN_STALE_REMOVABLE,
)
self.assertEqual(
wca.classify_worktree(
workflow_type=wca.WORKFLOW_CONFLICT_FIX,
is_dirty=False,
ttl_expired=False,
),
wca.CLASS_ACTIVE_ISSUE_WORK,
)
def test_issue_work_ttl_alone_no_longer_grants_removal(self):
"""Age is not landing proof: TTL alone must not reclaim issue work."""
self.assertEqual(
wca.classify_worktree(
workflow_type=wca.WORKFLOW_ISSUE_WORK,
is_dirty=False,
ttl_expired=True,
),
wca.CLASS_ACTIVE_ISSUE_WORK,
)
class TestAssessorPerformsNoDeletion(_AuditHarness):
def test_audit_never_removes_a_worktree(self):
with patch.object(wca, "remove_worktree") as removal:
self.merged_audit()
removal.assert_not_called()
def test_audit_shells_out_to_no_destructive_git_command(self):
seen = []
real_run = subprocess.run
def recording_run(cmd, *args, **kwargs):
seen.append(cmd)
return real_run(["true"], *args, **kwargs)
with patch.object(subprocess, "run", side_effect=recording_run):
wca.audit_branches_directory("/nonexistent-repo-for-audit")
joined = [" ".join(c) if isinstance(c, list) else str(c) for c in seen]
for cmd in joined:
self.assertNotIn("worktree remove", cmd)
self.assertNotIn("branch -D", cmd)
self.assertNotIn("push", cmd)
class TestAgreementWithPrScopedReconciler(unittest.TestCase):
"""The audit and merged_cleanup_reconcile must agree on identical input.
Uses a real throwaway git repository so containment is computed by git
rather than asserted. Nothing outside the temporary directory is touched.
"""
def _git(self, *args):
subprocess.run(
["git", "-C", self.root, *args],
check=True,
capture_output=True,
text=True,
)
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.root = os.path.realpath(self._tmp.name)
self._git("init", "-b", "master", ".")
self._git("config", "user.email", "[email protected]")
self._git("config", "user.name", "Test")
with open(os.path.join(self.root, "seed.txt"), "w") as fh:
fh.write("seed\n")
self._git("add", "seed.txt")
self._git("commit", "-m", "seed")
self.branch = "feat/issue-777-timeline"
self._git("checkout", "-b", self.branch)
with open(os.path.join(self.root, "feature.txt"), "w") as fh:
fh.write("feature\n")
self._git("add", "feature.txt")
self._git("commit", "-m", "feature")
self.head_sha = subprocess.run(
["git", "-C", self.root, "rev-parse", "HEAD"],
capture_output=True, text=True, check=True,
).stdout.strip()
self._git("checkout", "master")
self._git("merge", "--no-ff", "-m", "merge feature", self.branch)
self.worktree = os.path.join(self.root, "branches", "issue-777-timeline")
self._git("worktree", "add", self.worktree, self.branch)
def tearDown(self):
self._tmp.cleanup()
def _pr_index(self):
return wca.build_pr_index(
[
{
"number": 849,
"head": {"ref": self.branch, "sha": self.head_sha},
"merged_at": "2026-07-24T01:00:00Z",
}
]
)
def _audit_entry(self):
report = wca.audit_branches_directory(
self.root, pr_index=self._pr_index(), master_ref="master"
)
return next(wt for wt in report["worktrees"] if wt["path"] == self.worktree)
def _reconciler_entry(self):
return mcr.assess_local_worktree_cleanup(
pr_number=849,
head_branch=self.branch,
merged=True,
worktree_state=mcr.resolve_cleanup_worktree_state(
project_root=self.root,
head_branch=self.branch,
issue_number=777,
pr_head_sha=self.head_sha,
target_ref="master",
),
active_lock=False,
)
def test_both_assessors_agree_the_worktree_is_safe(self):
audit_entry = self._audit_entry()
reconciler = self._reconciler_entry()
self.assertTrue(reconciler["safe_to_remove_worktree"], reconciler)
self.assertTrue(audit_entry["removable"], audit_entry)
self.assertEqual(audit_entry["pr_number"], reconciler["pr_number"])
self.assertEqual(audit_entry["merged_pr_cleanup"]["block_reasons"], [])
self.assertEqual(reconciler["block_reasons"], [])
def test_both_assessors_agree_a_dirty_worktree_is_unsafe(self):
with open(os.path.join(self.worktree, "feature.txt"), "a") as fh:
fh.write("local edit\n")
audit_entry = self._audit_entry()
reconciler = self._reconciler_entry()
self.assertFalse(audit_entry["removable"])
self.assertFalse(reconciler["safe_to_remove_worktree"])
def test_worktree_still_present_after_audit(self):
self._audit_entry()
self.assertTrue(os.path.isdir(self.worktree))
if __name__ == "__main__":
unittest.main()
+24 -2
View File
@@ -134,13 +134,35 @@ class TestClassification(unittest.TestCase):
self.assertEqual(cls, wca.CLASS_ACTIVE_OPEN_PR) self.assertEqual(cls, wca.CLASS_ACTIVE_OPEN_PR)
self.assertFalse(wca.is_removable(cls)) self.assertFalse(wca.is_removable(cls))
def test_stale_clean_issue_worktree_removable(self): def test_stale_clean_issue_worktree_needs_merged_pr_proof(self):
# Scenario 5: clean issue worktree, TTL expired, no lock -> removable. # Scenario 5 (#858): age is not proof that the branch landed, so a
# TTL-expired issue worktree stays active work. Only authoritative
# merged-PR evidence makes it removable, which is what keeps a
# worktree holding unmerged commits from being reclaimed by age.
cls = wca.classify_worktree( cls = wca.classify_worktree(
workflow_type=wca.WORKFLOW_ISSUE_WORK, workflow_type=wca.WORKFLOW_ISSUE_WORK,
is_dirty=False, is_dirty=False,
ttl_expired=True, ttl_expired=True,
) )
self.assertEqual(cls, wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(wca.is_removable(cls))
cls = wca.classify_worktree(
workflow_type=wca.WORKFLOW_ISSUE_WORK,
is_dirty=False,
ttl_expired=True,
merged_pr_cleanup={"proven": True},
)
self.assertEqual(cls, wca.CLASS_CLEAN_STALE_REMOVABLE)
self.assertTrue(wca.is_removable(cls))
def test_stale_clean_conflict_fix_worktree_removable(self):
# conflict_fix keeps the original TTL rule; #858 changed issue work only.
cls = wca.classify_worktree(
workflow_type=wca.WORKFLOW_CONFLICT_FIX,
is_dirty=False,
ttl_expired=True,
)
self.assertEqual(cls, wca.CLASS_CLEAN_STALE_REMOVABLE) self.assertEqual(cls, wca.CLASS_CLEAN_STALE_REMOVABLE)
self.assertTrue(wca.is_removable(cls)) self.assertTrue(wca.is_removable(cls))
+271 -4
View File
@@ -34,7 +34,11 @@ import subprocess
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Any from typing import Any
from merged_cleanup_reconcile import branch_worktree_folder, read_local_worktree_state from merged_cleanup_reconcile import (
branch_worktree_folder,
is_head_ancestor_of_ref,
read_local_worktree_state,
)
from reviewer_worktree import parse_dirty_tracked_files, REVIEW_WORKTREE_RE from reviewer_worktree import parse_dirty_tracked_files, REVIEW_WORKTREE_RE
PROTECTED_BRANCHES = frozenset({"master", "main", "dev"}) PROTECTED_BRANCHES = frozenset({"master", "main", "dev"})
@@ -67,6 +71,14 @@ REMOVABLE_CLASSES = frozenset(
{CLASS_CLEAN_STALE_REMOVABLE, CLASS_DETACHED_REVIEW_LEFTOVER} {CLASS_CLEAN_STALE_REMOVABLE, CLASS_DETACHED_REVIEW_LEFTOVER}
) )
# Merged-PR linkage outcomes for issue worktrees (#858). Only ``LINKAGE_MERGED``
# is ownership proof; every other outcome leaves the worktree protected.
LINKAGE_MERGED = "merged_pr"
LINKAGE_OPEN = "open_pr"
LINKAGE_NONE = "no_owning_pr"
LINKAGE_AMBIGUOUS = "ambiguous"
LINKAGE_UNKNOWN = "unknown"
_ISSUE_REF_RE = re.compile(r"issue-(\d+)", re.IGNORECASE) _ISSUE_REF_RE = re.compile(r"issue-(\d+)", re.IGNORECASE)
_ISSUE_BRANCH_PREFIXES = ("feat/", "fix/", "docs/", "chore/") _ISSUE_BRANCH_PREFIXES = ("feat/", "fix/", "docs/", "chore/")
@@ -169,6 +181,186 @@ def is_ttl_expired(
return (now_dt - last).total_seconds() > ttl_hours * 3600.0 return (now_dt - last).total_seconds() > ttl_hours * 3600.0
def build_pr_index(prs: list[dict[str, Any]] | None) -> dict[str, list[dict[str, Any]]]:
"""Index PR records by head branch for deterministic worktree linkage (#858).
Accepts Gitea PR payloads (``head`` as a dict) and pre-flattened records
(``head_branch``/``head_sha``). Records without a usable head branch or
number are dropped rather than guessed at, so a branch is only ever linked
to a PR the caller actually proved.
"""
index: dict[str, list[dict[str, Any]]] = {}
for pr in prs or []:
head = pr.get("head")
if isinstance(head, dict):
head_branch = head.get("ref")
head_sha = head.get("sha")
else:
head_branch = pr.get("head_branch") or (head if isinstance(head, str) else None)
head_sha = pr.get("head_sha")
number = pr.get("number")
if not head_branch or number is None:
continue
try:
pr_number = int(number)
except (TypeError, ValueError):
continue
index.setdefault(str(head_branch).strip(), []).append(
{
"pr_number": pr_number,
"head_branch": str(head_branch).strip(),
"head_sha": head_sha,
"merged": bool(pr.get("merged") or pr.get("merged_at")),
"state": pr.get("state"),
}
)
return index
def resolve_owning_pr(
*,
branch: str | None,
pr_index: dict[str, list[dict[str, Any]]] | None,
) -> dict[str, Any]:
"""Resolve the single PR that owns ``branch``, failing closed when unclear.
Ownership is only ``LINKAGE_MERGED`` when exactly one PR claims the branch
and that PR is merged. Several distinct PRs on one branch is a competing
claim (``LINKAGE_AMBIGUOUS``), and a still-open owner is reported as
``LINKAGE_OPEN`` — both keep the worktree protected while still exposing
the PR number the audit resolved.
"""
if pr_index is None:
return {
"status": LINKAGE_UNKNOWN,
"pr_number": None,
"candidate_pr_numbers": [],
"reasons": ["live PR state was not supplied; ownership unproven"],
}
branch_name = (branch or "").strip()
if not branch_name:
return {
"status": LINKAGE_UNKNOWN,
"pr_number": None,
"candidate_pr_numbers": [],
"reasons": ["worktree has no attached branch; ownership unproven"],
}
candidates = list(pr_index.get(branch_name) or [])
numbers = sorted({c["pr_number"] for c in candidates})
if not candidates:
return {
"status": LINKAGE_NONE,
"pr_number": None,
"candidate_pr_numbers": [],
"reasons": [f"no PR claims branch '{branch_name}'"],
}
if len(numbers) > 1:
return {
"status": LINKAGE_AMBIGUOUS,
"pr_number": None,
"candidate_pr_numbers": numbers,
"reasons": [
f"branch '{branch_name}' is claimed by competing PRs {numbers}; "
"ownership is ambiguous"
],
}
owner = candidates[0]
pr_number = owner["pr_number"]
if owner.get("head_branch") != branch_name:
return {
"status": LINKAGE_UNKNOWN,
"pr_number": pr_number,
"candidate_pr_numbers": numbers,
"reasons": [
f"PR #{pr_number} head branch '{owner.get('head_branch')}' does not "
f"match worktree branch '{branch_name}'"
],
}
if not owner.get("merged"):
return {
"status": LINKAGE_OPEN,
"pr_number": pr_number,
"candidate_pr_numbers": numbers,
"pr_head_sha": owner.get("head_sha"),
"reasons": [f"owning PR #{pr_number} is not merged"],
}
return {
"status": LINKAGE_MERGED,
"pr_number": pr_number,
"candidate_pr_numbers": numbers,
"pr_head_sha": owner.get("head_sha"),
"reasons": [],
}
def assess_merged_pr_worktree_cleanup(
*,
linkage: dict[str, Any] | None,
head_sha: str | None,
head_in_master: bool | None,
is_dirty: bool,
has_open_pr: bool,
has_active_lease: bool,
has_active_issue_lock: bool,
is_protected: bool,
has_live_session: bool = False,
) -> dict[str, Any]:
"""Decide whether a merged issue worktree satisfies the full cleanup policy.
Every condition must be independently proven: conclusive merged-PR
ownership, agreement between the worktree branch and the PR head branch,
containment of the worktree head in authoritative master (which is what
proves no unmerged commits remain), absence of any open/competing PR,
lease, issue lock, or live session, a clean tree, and a worktree that is
not the protected control checkout. Anything unknown blocks.
"""
link = linkage or {
"status": LINKAGE_UNKNOWN,
"pr_number": None,
"reasons": ["no linkage assessment supplied"],
}
status = link.get("status")
reasons: list[str] = []
if status != LINKAGE_MERGED:
reasons.extend(
link.get("reasons") or ["owning PR could not be conclusively identified"]
)
if is_protected:
reasons.append("worktree is protected or the stable control checkout")
if is_dirty:
reasons.append("worktree has uncommitted changes")
if has_open_pr:
reasons.append("worktree branch has an open PR")
if has_active_lease:
reasons.append("worktree has an active lease")
if has_active_issue_lock:
reasons.append("an active issue lock references this branch")
if has_live_session:
reasons.append("a live process or session is using this worktree")
if not head_sha:
reasons.append("worktree head sha is unknown")
if head_in_master is None:
reasons.append("containment of the worktree head in master is unknown")
elif not head_in_master:
reasons.append(
"worktree head is not contained in authoritative master "
"(unmerged commits remain)"
)
proven = not reasons
return {
"linkage_status": status,
"pr_number": link.get("pr_number"),
"pr_head_sha": link.get("pr_head_sha"),
"head_in_master": head_in_master,
"proven": proven,
"block_reasons": reasons,
}
def classify_worktree( def classify_worktree(
*, *,
workflow_type: str, workflow_type: str,
@@ -181,6 +373,8 @@ def classify_worktree(
ttl_expired: bool = False, ttl_expired: bool = False,
is_protected: bool = False, is_protected: bool = False,
metadata_known: bool = True, metadata_known: bool = True,
merged_pr_cleanup: dict[str, Any] | None = None,
has_live_session: bool = False,
) -> str: ) -> str:
"""Classify a worktree, safety-first: any preservation signal wins. """Classify a worktree, safety-first: any preservation signal wins.
@@ -199,6 +393,8 @@ def classify_worktree(
return CLASS_ACTIVE_ISSUE_WORK # never auto-deleted (criterion 8) return CLASS_ACTIVE_ISSUE_WORK # never auto-deleted (criterion 8)
if has_active_issue_lock: if has_active_issue_lock:
return CLASS_ACTIVE_ISSUE_WORK return CLASS_ACTIVE_ISSUE_WORK
if has_live_session:
return CLASS_ACTIVE_ISSUE_WORK # a live session still owns this tree
if not metadata_known or workflow_type == WORKFLOW_UNKNOWN: if not metadata_known or workflow_type == WORKFLOW_UNKNOWN:
return CLASS_UNSAFE_UNKNOWN # never auto-deleted without proof return CLASS_UNSAFE_UNKNOWN # never auto-deleted without proof
@@ -207,7 +403,15 @@ def classify_worktree(
if is_detached or branch_gone: if is_detached or branch_gone:
return CLASS_DETACHED_REVIEW_LEFTOVER return CLASS_DETACHED_REVIEW_LEFTOVER
return CLASS_CLEAN_STALE_REMOVABLE return CLASS_CLEAN_STALE_REMOVABLE
# issue_work / conflict_fix: only removable once the TTL has expired. if workflow_type == WORKFLOW_ISSUE_WORK:
# #858: an issue worktree becomes removable only on authoritative
# merged-PR evidence satisfying the whole cleanup policy. Age alone
# never proves the branch landed, so TTL cannot qualify one by itself
# — otherwise a worktree holding unmerged commits would be reclaimed.
if (merged_pr_cleanup or {}).get("proven"):
return CLASS_CLEAN_STALE_REMOVABLE
return CLASS_ACTIVE_ISSUE_WORK
# conflict_fix: only removable once the TTL has expired.
if ttl_expired: if ttl_expired:
return CLASS_CLEAN_STALE_REMOVABLE return CLASS_CLEAN_STALE_REMOVABLE
return CLASS_ACTIVE_ISSUE_WORK return CLASS_ACTIVE_ISSUE_WORK
@@ -400,6 +604,20 @@ def remove_worktree(project_root: str, path: str) -> dict[str, Any]:
} }
def head_contained_in_ref(
project_root: str, head_sha: str | None, ref: str | None
) -> bool | None:
"""Return True when ``head_sha`` is already contained in ``ref``.
Shares :mod:`merged_cleanup_reconcile`'s ancestry check so the audit and
the PR-scoped reconciler agree on what "already landed" means (#858).
Returns None when containment cannot be determined, which fails closed.
"""
if not head_sha or not ref:
return None
return is_head_ancestor_of_ref(project_root, head_sha, ref)
def _is_under_branches(project_root: str, path: str) -> bool: def _is_under_branches(project_root: str, path: str) -> bool:
branches_root = os.path.join(os.path.abspath(project_root), "branches") branches_root = os.path.join(os.path.abspath(project_root), "branches")
return os.path.abspath(path or "").startswith(branches_root + os.sep) return os.path.abspath(path or "").startswith(branches_root + os.sep)
@@ -413,16 +631,30 @@ def audit_branches_directory(
active_issue_branches: set[str] | None = None, active_issue_branches: set[str] | None = None,
now: datetime | str | None = None, now: datetime | str | None = None,
ttl_hours: float = DEFAULT_TTL_HOURS, ttl_hours: float = DEFAULT_TTL_HOURS,
pr_index: dict[str, list[dict[str, Any]]] | None = None,
leased_issue_numbers: set[int] | None = None,
live_session_paths: set[str] | None = None,
master_ref: str | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Classify every session-owned worktree under ``branches/``. """Classify every session-owned worktree under ``branches/``.
Read-only: shells out to git for discovery and dirty state, then applies Read-only: shells out to git for discovery and dirty state, then applies
the pure classifier. Returns per-worktree classifications, counts, the the pure classifier. Returns per-worktree classifications, counts, the
list of removable candidates, and the ``git worktree list`` proof. list of removable candidates, and the ``git worktree list`` proof.
``pr_index`` (see :func:`build_pr_index`) supplies the authoritative PR
ownership used to link issue worktrees to their merged PR (#858).
``master_ref`` is the ref a worktree head must be contained in before it
can be considered landed. Both are optional and their absence only ever
fails closed: without them no issue worktree becomes removable.
""" """
open_pr_branches = open_pr_branches or set() open_pr_branches = open_pr_branches or set()
leased_branches = leased_branches or set() leased_branches = leased_branches or set()
active_issue_branches = active_issue_branches or set() active_issue_branches = active_issue_branches or set()
leased_issue_numbers = leased_issue_numbers or set()
live_session_paths = {
os.path.abspath(p) for p in (live_session_paths or set()) if p
}
worktrees: list[dict[str, Any]] = [] worktrees: list[dict[str, Any]] = []
for entry in list_worktrees(project_root): for entry in list_worktrees(project_root):
@@ -433,12 +665,42 @@ def audit_branches_directory(
) )
dirty_state = read_worktree_dirty(path) dirty_state = read_worktree_dirty(path)
is_dirty = bool(dirty_state.get("dirty")) is_dirty = bool(dirty_state.get("dirty"))
head_sha = entry.get("head")
linkage = resolve_owning_pr(branch=branch, pr_index=pr_index)
metadata = build_worktree_metadata( metadata = build_worktree_metadata(
path=path, branch=branch, head_sha=entry.get("head") path=path,
branch=branch,
head_sha=head_sha,
pr_number=linkage.get("pr_number"),
) )
has_open_pr = bool(branch) and branch in open_pr_branches has_open_pr = bool(branch) and branch in open_pr_branches
has_active_lease = bool(branch) and branch in leased_branches # A lease on issue N protects that issue's own work worktree. It must
# not incidentally protect a baseline/review scratch tree that merely
# carries the same issue marker in its name, which would change the
# classification of worktrees this policy does not own.
has_active_lease = (bool(branch) and branch in leased_branches) or (
metadata["workflow_type"] == WORKFLOW_ISSUE_WORK
and metadata.get("issue_number") is not None
and metadata["issue_number"] in leased_issue_numbers
)
has_active_lock = bool(branch) and branch in active_issue_branches has_active_lock = bool(branch) and branch in active_issue_branches
has_live_session = bool(path) and os.path.abspath(path) in live_session_paths
head_in_master = (
head_contained_in_ref(project_root, head_sha, master_ref)
if master_ref
else None
)
merged_pr_cleanup = assess_merged_pr_worktree_cleanup(
linkage=linkage,
head_sha=head_sha,
head_in_master=head_in_master,
is_dirty=is_dirty,
has_open_pr=has_open_pr,
has_active_lease=has_active_lease,
has_active_issue_lock=has_active_lock,
is_protected=is_protected,
has_live_session=has_live_session,
)
ttl_expired = is_ttl_expired( ttl_expired = is_ttl_expired(
last_used_at=metadata.get("last_used_at"), now=now, ttl_hours=ttl_hours last_used_at=metadata.get("last_used_at"), now=now, ttl_hours=ttl_hours
) )
@@ -452,6 +714,8 @@ def audit_branches_directory(
branch_gone=branch is None and not entry.get("detached"), branch_gone=branch is None and not entry.get("detached"),
ttl_expired=ttl_expired, ttl_expired=ttl_expired,
is_protected=is_protected, is_protected=is_protected,
merged_pr_cleanup=merged_pr_cleanup,
has_live_session=has_live_session,
) )
metadata["cleanup_eligibility"] = classification metadata["cleanup_eligibility"] = classification
worktrees.append( worktrees.append(
@@ -463,7 +727,10 @@ def audit_branches_directory(
"has_open_pr": has_open_pr, "has_open_pr": has_open_pr,
"has_active_lease": has_active_lease, "has_active_lease": has_active_lease,
"has_active_issue_lock": has_active_lock, "has_active_issue_lock": has_active_lock,
"has_live_session": has_live_session,
"is_protected": is_protected, "is_protected": is_protected,
"merged_pr_linkage": linkage,
"merged_pr_cleanup": merged_pr_cleanup,
"classification": classification, "classification": classification,
"removable": is_removable(classification), "removable": is_removable(classification),
} }