From 18d6583e8362b15621b0a9f4766c6e073419e627 Mon Sep 17 00:00:00 2001 From: jcwalker3 Date: Thu, 23 Jul 2026 20:38:47 -0500 Subject: [PATCH 1/3] fix(author): bootstrap recovery for dirty orphaned issue worktrees (#860) Add an explicit recovery operation for same-claimant dirty registered worktrees under malformed PID-less durable locks, with crash-safe journals, dirty byte preservation, path-level conflict detection, and live session binding. PID-less locks are never treated as live merely because expiry is absent. Closes #860 Co-Authored-By: Grok 4.5 (xAI) --- dirty_orphan_worktree_recovery.py | 1069 ++++++++++++++++++ gitea_mcp_server.py | 242 ++++ issue_lock_store.py | 37 +- task_capability_map.py | 9 + tests/test_dirty_orphan_worktree_recovery.py | 422 +++++++ tests/test_issue_lock_store.py | 6 + 6 files changed, 1782 insertions(+), 3 deletions(-) create mode 100644 dirty_orphan_worktree_recovery.py create mode 100644 tests/test_dirty_orphan_worktree_recovery.py diff --git a/dirty_orphan_worktree_recovery.py b/dirty_orphan_worktree_recovery.py new file mode 100644 index 0000000..fdbcfa6 --- /dev/null +++ b/dirty_orphan_worktree_recovery.py @@ -0,0 +1,1069 @@ +"""Dirty orphaned author-issue worktree recovery (#860). + +Self-hosting deadlock class observed against #850 / PR #853 and #855: + +* same-claimant durable issue lock (``jcwalker3`` / ``prgs-author``) +* registered dirty worktree under ``branches/`` +* lock missing PID / session PID / expiry / heartbeat (malformed) +* existing recovery (#753/#768/#772) and renewal require a *clean* worktree + and/or a determinable owner PID +* no sanctioned dirty-preserving rebind + sync to a newer remote PR head + +This module is the pure evidence assessor and crash-safe recovery orchestrator +for that one class. It is an **explicit** recovery operation — it does not +silently widen ``gitea_lock_issue``. + +Safety model +------------ +Eligibility requires *all* of: + +* same claimant identity and profile as the durable lock +* exact issue / repository / branch / registered source worktree agreement +* registered worktree under the canonical branches root (resolved-path ancestry) +* no active original process (when PID is present) +* no active competing author session or workflow lease +* sufficient corroborating evidence when PID fields are absent (caller pins) +* explicit caller-provided local head, remote/PR head, and dirty fingerprints +* no foreign, duplicate, or ambiguous ownership + +A PID-less lock is **never** considered live merely because expiration fields +are absent (see also ``issue_lock_store.assess_lock_freshness``). + +Dirty preservation +------------------ +The source worktree is frozen: never cleaned, reset, overwritten, or deleted. +Recovery prefers a separately prepared recovery worktree checked out at the +pinned remote PR head. Dirty bytes are re-applied with path-level conflict +detection when upstream also changed a dirty path. + +Crash safety +------------ +A durable journal is written *before* filesystem or ownership mutation. +Retries resume or fail closed without stealing ownership, duplicating +worktrees, or losing dirty bytes. +""" + +from __future__ import annotations + +import hashlib +import json +import os +import shutil +import stat +import subprocess +from dataclasses import dataclass +from typing import Any, Mapping, Sequence + +from issue_lock_store import is_process_alive +from reviewer_worktree import parse_dirty_tracked_files + +# Outcomes +ELIGIBLE = "ELIGIBLE" +REFUSED = "REFUSED" +NO_CANDIDATE = "NO_CANDIDATE" +RECOVERY_COMPLETED = "RECOVERY_COMPLETED" +RECOVERY_RESUMED = "RECOVERY_RESUMED" +CONFLICTS_PRESENT = "CONFLICTS_PRESENT" + +# Journal phases (ordered) +PHASE_ELIGIBILITY = "1_eligibility_proven" +PHASE_JOURNAL_PERSISTED = "2_journal_persisted" +PHASE_RECOVERY_WORKTREE = "3_recovery_worktree_prepared" +PHASE_DIRTY_APPLIED = "4_dirty_applied" +PHASE_BINDING = "5_session_bound" +PHASE_COMPLETE = "6_complete" + +JOURNAL_DIR_NAME = "dirty-orphan-recovery-journals" +REQUIRED_LOCK_FIELDS = ("issue_number", "branch_name", "worktree_path") + +# Marker directory written into recovery worktrees for governed conflicts. +CONFLICT_STATE_DIR = ".gitea-recovery" +CONFLICT_STATE_FILE = "conflicts.json" + + +def _text(value: Any) -> str: + return str(value or "").strip() + + +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 sha256_bytes(data: bytes) -> str: + return hashlib.sha256(data).hexdigest() + + +def sha256_file(path: str) -> str: + with open(path, "rb") as fh: + return sha256_bytes(fh.read()) + + +def get_journal_dir(override: str | None = None) -> str: + if override: + path = override + elif os.environ.get("GITEA_DIRTY_ORPHAN_RECOVERY_JOURNAL_DIR"): + path = os.environ["GITEA_DIRTY_ORPHAN_RECOVERY_JOURNAL_DIR"] + else: + path = os.path.join( + os.path.expanduser("~/.cache/gitea-tools"), JOURNAL_DIR_NAME + ) + os.makedirs(path, mode=0o700, exist_ok=True) + return path + + +def _journal_path(idempotency_key: str, journal_dir: str | None = None) -> str: + safe = "".join( + c if c.isalnum() or c in ("-", "_", ".") else "_" + for c in idempotency_key + ) + return os.path.join(get_journal_dir(journal_dir), f"{safe}.json") + + +def load_journal( + idempotency_key: str, journal_dir: str | None = None +) -> dict[str, Any] | None: + path = _journal_path(idempotency_key, journal_dir=journal_dir) + if not os.path.isfile(path): + return None + if os.path.islink(path): + raise ValueError(f"refusing journal path that is a symlink: {path}") + with open(path, "r", encoding="utf-8") as fh: + data = json.load(fh) + if not isinstance(data, dict): + raise ValueError("corrupt recovery journal: not an object") + return data + + +def save_journal( + journal: Mapping[str, Any], journal_dir: str | None = None +) -> str: + key = _text(journal.get("idempotency_key")) + if not key: + raise ValueError("journal requires idempotency_key") + path = _journal_path(key, journal_dir=journal_dir) + if os.path.islink(path): + raise ValueError(f"refusing to write journal through symlink: {path}") + tmp = f"{path}.tmp.{os.getpid()}" + with open(tmp, "w", encoding="utf-8") as fh: + json.dump(dict(journal), fh, indent=2, sort_keys=True) + fh.flush() + os.fsync(fh.fileno()) + os.replace(tmp, path) + return path + + +def derive_idempotency_key( + *, + remote: str, + org: str, + repo: str, + issue_number: int, + source_worktree: str, + expected_local_head: str, + expected_remote_head: str, +) -> str: + raw = "|".join( + [ + "dirty-orphan-recovery", + _text(remote), + _text(org), + _text(repo), + f"issue-{int(issue_number)}", + os.path.realpath(_text(source_worktree)), + _text(expected_local_head)[:40], + _text(expected_remote_head)[:40], + ] + ) + return hashlib.sha256(raw.encode("utf-8")).hexdigest()[:32] + + +def _lock_claimant(lock: Mapping[str, Any]) -> dict[str, str]: + claimant = lock.get("claimant") + if not isinstance(claimant, Mapping): + lease = lock.get("work_lease") + claimant = lease.get("claimant") if isinstance(lease, Mapping) else None + if not isinstance(claimant, Mapping): + return {} + return { + "username": _text(claimant.get("username")), + "profile": _text(claimant.get("profile")), + } + + +def _recorded_pid(lock: Mapping[str, Any]) -> Any: + pid = lock.get("session_pid") + if pid is None: + pid = lock.get("pid") + return pid + + +def _malformed_pid_fields(lock: Mapping[str, Any]) -> list[str]: + """Return reasons the lock's PID fields are unusable (not automatically live).""" + missing: list[str] = [] + pid = _recorded_pid(lock) + if pid is None or _text(pid) == "": + missing.append("session_pid/pid") + return missing + try: + if int(pid) <= 0: + missing.append("session_pid/pid") + except (TypeError, ValueError): + missing.append("session_pid/pid") + return missing + + +def is_path_under_canonical_branches( + path: str, + *, + canonical_repo_root: str, + branches_dirname: str = "branches", +) -> tuple[bool, list[str]]: + """Resolved-path ancestry under ``{canonical_repo_root}/branches``. + + Rejects string-substring tricks, traversal, and symlink escapes outside + the canonical branches root. + """ + reasons: list[str] = [] + if not path or not canonical_repo_root: + return False, ["path and canonical_repo_root are required"] + try: + repo_root = os.path.realpath(canonical_repo_root) + branches_root = os.path.realpath(os.path.join(repo_root, branches_dirname)) + # realpath on a non-existent path still normalizes; prefer it so + # assessment can run without the directory existing yet. + target = os.path.realpath(path) + except OSError as exc: + return False, [f"path resolution failed: {exc}"] + + if not target.startswith(branches_root + os.sep) and target != branches_root: + reasons.append( + f"path '{path}' is not under canonical branches root '{branches_root}'" + ) + return False, reasons + try: + rel = os.path.relpath(target, branches_root) + except ValueError: + return False, ["path not relative to branches root"] + if rel.startswith(".."): + return False, ["path escapes branches root via relative traversal"] + return True, [] + + +def _result( + outcome: str, + *, + eligible: bool, + reasons: list[str], + evidence: dict[str, Any], +) -> dict[str, Any]: + return { + "outcome": outcome, + "eligible": eligible, + "recovery_eligible": eligible, + "reasons": list(reasons), + "evidence": evidence, + } + + +def assess_dirty_orphan_recovery( + existing_lock: Mapping[str, Any] | None, + *, + issue_number: int, + branch_name: str, + source_worktree_path: str, + remote: str, + org: str, + repo: str, + identity: str, + profile: str, + expected_local_head: str, + expected_remote_head: str, + expected_dirty_fingerprints: Mapping[str, str], + current_branch: str | None, + porcelain_status: str, + observed_local_head: str | None, + observed_remote_head: str | None, + observed_dirty_fingerprints: Mapping[str, str] | None, + competing_live_locks: Sequence[Mapping[str, Any]] | None = None, + competing_live_sessions: Sequence[Mapping[str, Any]] | None = None, + workflow_lease_active: bool | None = None, + workflow_lease_expired: bool | None = None, + canonical_repo_root: str, + worktree_registered: bool, + current_pid: int | None = None, + owner_process_alive_override: bool | None = None, +) -> dict[str, Any]: + """Assess whether dirty-orphan recovery is eligible. Pure; no I/O mutation.""" + evidence: dict[str, Any] = { + "issue_number": issue_number, + "branch_name": branch_name, + "source_worktree_path": source_worktree_path, + "remote": remote, + "org": org, + "repo": repo, + "identity": identity, + "profile": profile, + "expected_local_head": _text(expected_local_head), + "expected_remote_head": _text(expected_remote_head), + "expected_dirty_paths": sorted(expected_dirty_fingerprints or {}), + } + reasons: list[str] = [] + + if not existing_lock: + return _result( + NO_CANDIDATE, + eligible=False, + reasons=["no existing durable lock for this issue"], + evidence=evidence, + ) + + lock = dict(existing_lock) + if lock.get("issue_number") != issue_number: + return _result( + NO_CANDIDATE, + eligible=False, + reasons=[ + f"existing lock targets issue #{lock.get('issue_number')}, " + f"not #{issue_number}" + ], + evidence=evidence, + ) + + for field in REQUIRED_LOCK_FIELDS: + if not _text(lock.get(field)): + reasons.append(f"durable lock missing required field '{field}'") + + # Repository / branch / worktree agreement + for field, expected in (("remote", remote), ("org", org), ("repo", repo)): + actual = _text(lock.get(field)) + if actual and actual != _text(expected): + reasons.append( + f"lock {field} '{actual}' does not match requested '{_text(expected)}'" + ) + elif not actual: + # Some legacy locks omit remote/org/repo; require explicit pin match + # via caller still supplying them and branch/worktree agreement. + evidence[f"lock_{field}_absent"] = True + + locked_branch = _text(lock.get("branch_name")) + if locked_branch != _text(branch_name): + reasons.append( + f"lock branch '{locked_branch}' does not match requested " + f"'{_text(branch_name)}'" + ) + evidence["locked_branch"] = locked_branch + + locked_wt = _text(lock.get("worktree_path")) + if not _same_realpath(locked_wt, source_worktree_path): + reasons.append( + f"lock worktree '{locked_wt}' does not match declared source " + f"'{_text(source_worktree_path)}'" + ) + evidence["locked_worktree_path"] = locked_wt + + under, under_reasons = is_path_under_canonical_branches( + source_worktree_path, canonical_repo_root=canonical_repo_root + ) + if not under: + reasons.extend(under_reasons) + if not worktree_registered: + reasons.append("source worktree is not registered in git worktree list") + + # Claimant identity + claimant = _lock_claimant(lock) + evidence["lock_claimant"] = claimant + if claimant.get("username") != _text(identity): + reasons.append( + f"foreign claimant identity '{claimant.get('username')}' " + f"(caller '{_text(identity)}')" + ) + if claimant.get("profile") != _text(profile): + reasons.append( + f"foreign claimant profile '{claimant.get('profile')}' " + f"(caller '{_text(profile)}')" + ) + + # PID / liveness + pid_missing = _malformed_pid_fields(lock) + recorded_pid = _recorded_pid(lock) + evidence["recorded_pid"] = recorded_pid + evidence["pid_fields_missing"] = pid_missing + if pid_missing: + evidence["pid_less_malformed"] = True + # PID-less is never live by missing expiry alone. Eligibility continues + # only with full corroborating pins (already required below). + else: + alive = ( + owner_process_alive_override + if owner_process_alive_override is not None + else is_process_alive(int(recorded_pid)) + ) + evidence["owner_process_alive"] = alive + if alive: + reasons.append( + f"original owner process pid {recorded_pid} is still alive; " + "recovery refused" + ) + + # Competing ownership + for entry in competing_live_locks or (): + if not isinstance(entry, Mapping): + continue + if entry.get("issue_number") == issue_number: + reasons.append( + "competing live lock observed for the same issue; recovery refused" + ) + for entry in competing_live_sessions or (): + if not isinstance(entry, Mapping): + continue + sess_issue = entry.get("issue_number") + sess_identity = _text(entry.get("identity") or entry.get("username")) + if sess_issue == issue_number and sess_identity and sess_identity != _text(identity): + reasons.append( + f"competing live author session by '{sess_identity}' on issue " + f"#{issue_number}" + ) + elif sess_issue == issue_number and entry.get("active"): + # same claimant active elsewhere still blocks ambiguous ownership + if entry.get("session_pid") not in (None, current_pid, os.getpid()): + reasons.append( + "ambiguous competing same-issue author session still active" + ) + + if workflow_lease_active is True and workflow_lease_expired is not True: + reasons.append( + "active competing workflow lease still live; recovery refused" + ) + evidence["workflow_lease_active"] = workflow_lease_active + evidence["workflow_lease_expired"] = workflow_lease_expired + + # Dirty requirement (this recovery class is *for* dirty trees) + dirty_files = parse_dirty_tracked_files(porcelain_status) + if not dirty_files and not expected_dirty_fingerprints: + reasons.append( + "worktree is clean and no dirty fingerprints were provided; " + "use clean-worktree recovery (#753/#772) instead" + ) + evidence["observed_dirty_files"] = dirty_files + + # Explicit pins + if not _text(expected_local_head) or len(_text(expected_local_head)) < 40: + reasons.append("expected_local_head pin missing or not a full SHA") + if not _text(expected_remote_head) or len(_text(expected_remote_head)) < 40: + reasons.append("expected_remote_head pin missing or not a full SHA") + if not expected_dirty_fingerprints: + reasons.append("expected_dirty_fingerprints pin is required") + + if _text(observed_local_head) and _text(observed_local_head) != _text( + expected_local_head + ): + reasons.append( + f"local head mismatch: observed {_text(observed_local_head)} != " + f"pinned {_text(expected_local_head)}" + ) + if _text(observed_remote_head) and _text(observed_remote_head) != _text( + expected_remote_head + ): + reasons.append( + f"remote/PR head mismatch: observed {_text(observed_remote_head)} != " + f"pinned {_text(expected_remote_head)}" + ) + if _text(expected_local_head) == _text(expected_remote_head): + # Divergence is the motivating case; equal heads are allowed only when + # dirty files still need rebinding, so do not refuse equality. + evidence["heads_equal"] = True + else: + evidence["heads_diverged"] = True + + if current_branch and _text(current_branch) != locked_branch: + reasons.append( + f"source worktree is on branch '{_text(current_branch)}', not " + f"locked branch '{locked_branch}'" + ) + + # Fingerprint verification + observed_fps = dict(observed_dirty_fingerprints or {}) + for path, expected_fp in (expected_dirty_fingerprints or {}).items(): + rel = _text(path) + if not rel or rel.startswith("/") or ".." in rel.split("/"): + reasons.append(f"unsafe dirty path pin refused: {path!r}") + continue + obs = _text(observed_fps.get(rel)) + if not obs: + reasons.append(f"missing observed fingerprint for dirty path '{rel}'") + elif obs != _text(expected_fp): + reasons.append( + f"dirty fingerprint mismatch for '{rel}': " + f"observed {obs} != pinned {_text(expected_fp)}" + ) + + # PID-less corroboration: all explicit pins must already have passed. + if pid_missing and reasons: + reasons.append( + "PID-less malformed lock additionally requires full pin corroboration; " + "one or more corroborating checks failed" + ) + + if reasons: + return _result(REFUSED, eligible=False, reasons=reasons, evidence=evidence) + + evidence["eligibility"] = ELIGIBLE + return _result(ELIGIBLE, eligible=True, reasons=[], evidence=evidence) + + +def detect_path_conflicts( + *, + dirty_paths: Sequence[str], + local_head_contents: Mapping[str, bytes | None], + remote_head_contents: Mapping[str, bytes | None], + dirty_contents: Mapping[str, bytes], +) -> list[dict[str, Any]]: + """Path-level conflicts: upstream and dirty patch both changed the path.""" + conflicts: list[dict[str, Any]] = [] + for path in dirty_paths: + local_b = local_head_contents.get(path) + remote_b = remote_head_contents.get(path) + dirty_b = dirty_contents.get(path) + if dirty_b is None: + continue + # Upstream changed relative to the local head version of the path. + upstream_changed = (local_b or b"") != (remote_b or b"") + dirty_differs_from_remote = dirty_b != (remote_b or b"") + if upstream_changed and dirty_differs_from_remote: + conflicts.append( + { + "path": path, + "local_head_sha256": sha256_bytes(local_b) if local_b is not None else None, + "remote_head_sha256": sha256_bytes(remote_b) if remote_b is not None else None, + "dirty_sha256": sha256_bytes(dirty_b), + "reason": ( + "upstream and preserved dirty patch both changed this path" + ), + } + ) + return conflicts + + +def _safe_open_lockfile(path: str): + """Open a lock file refusing symlinks (O_NOFOLLOW when available).""" + if os.path.islink(path): + raise ValueError(f"refusing lock file that is a symlink: {path}") + flags = os.O_RDWR | os.O_CREAT + if hasattr(os, "O_NOFOLLOW"): + flags |= os.O_NOFOLLOW + fd = os.open(path, flags, 0o600) + try: + st = os.fstat(fd) + if stat.S_ISLNK(st.st_mode): + os.close(fd) + raise ValueError(f"refusing lock file that is a symlink: {path}") + except Exception: + try: + os.close(fd) + except OSError: + pass + raise + return fd + + +def build_recovery_lock_record( + *, + existing_lock: Mapping[str, Any], + issue_number: int, + branch_name: str, + recovery_worktree_path: str, + remote: str, + org: str, + repo: str, + identity: str, + profile: str, + expected_remote_head: str, + source_worktree_path: str, + conflicts: Sequence[Mapping[str, Any]], + session_pid: int, +) -> dict[str, Any]: + """Construct a new durable lock bound to the recovery worktree + live PID.""" + from datetime import datetime, timedelta, timezone + + now = datetime.now(timezone.utc) + expires = now + timedelta(hours=4) + record = { + "remote": remote, + "org": org, + "repo": repo, + "issue_number": issue_number, + "branch_name": branch_name, + "worktree_path": recovery_worktree_path, + "pid": session_pid, + "session_pid": session_pid, + "claimant": {"username": identity, "profile": profile}, + "work_lease": { + "operation_type": "author_issue_work", + "issue_number": issue_number, + "pr_number": existing_lock.get("work_lease", {}).get("pr_number") + if isinstance(existing_lock.get("work_lease"), Mapping) + else existing_lock.get("pr_number"), + "branch": branch_name, + "worktree_path": recovery_worktree_path, + "claimant": {"username": identity, "profile": profile}, + "created_at": now.strftime("%Y-%m-%dT%H:%M:%SZ"), + "expires_at": expires.strftime("%Y-%m-%dT%H:%M:%SZ"), + "last_heartbeat_at": now.strftime("%Y-%m-%dT%H:%M:%SZ"), + }, + "dirty_orphan_recovery": { + "recovered": True, + "source_worktree_path": source_worktree_path, + "recovery_worktree_path": recovery_worktree_path, + "accepted_head": expected_remote_head, + "conflicts": list(conflicts), + "source_frozen": True, + }, + "lock_provenance": { + "source": "gitea_recover_dirty_orphaned_issue_worktree", + "written_by_tool": "gitea_recover_dirty_orphaned_issue_worktree", + "written_at": now.strftime("%Y-%m-%dT%H:%M:%SZ"), + "claimant": {"username": identity, "profile": profile}, + }, + } + # Preserve prior assignment/lease ids when present (non-authoritative). + for key in ("assignment_id", "lease_id", "owner_session", "expected_base_sha"): + if key in existing_lock: + record[key] = existing_lock[key] + return record + + +def write_conflict_state( + recovery_worktree: str, conflicts: Sequence[Mapping[str, Any]] +) -> str: + state_dir = os.path.join(recovery_worktree, CONFLICT_STATE_DIR) + os.makedirs(state_dir, mode=0o700, exist_ok=True) + path = os.path.join(state_dir, CONFLICT_STATE_FILE) + payload = { + "conflicts": list(conflicts), + "resolution": "author_edit_required" if conflicts else "none", + } + with open(path, "w", encoding="utf-8") as fh: + json.dump(payload, fh, indent=2, sort_keys=True) + return path + + +def apply_dirty_bytes( + *, + recovery_worktree: str, + dirty_contents: Mapping[str, bytes], + conflict_paths: set[str], +) -> list[str]: + """Write non-conflicting dirty bytes into the recovery worktree. + + Conflicting paths are written as ``*.recovered-dirty`` siblings so the + original dirty bytes remain recoverable without overwriting upstream. + """ + written: list[str] = [] + for rel, data in dirty_contents.items(): + rel = _text(rel) + if not rel or rel.startswith("/") or ".." in rel.split("/"): + raise ValueError(f"unsafe relative path: {rel!r}") + dest = os.path.join(recovery_worktree, rel) + parent = os.path.dirname(dest) + if parent: + os.makedirs(parent, exist_ok=True) + if rel in conflict_paths: + sidecar = dest + ".recovered-dirty" + with open(sidecar, "wb") as fh: + fh.write(data) + written.append(rel + ".recovered-dirty") + else: + with open(dest, "wb") as fh: + fh.write(data) + written.append(rel) + return written + + +@dataclass +class GitOps: + """Injectable git operations for tests.""" + + def run(self, args: list[str], *, cwd: str) -> subprocess.CompletedProcess[str]: + return subprocess.run( + args, + cwd=cwd, + capture_output=True, + text=True, + check=False, + ) + + +def prepare_recovery_worktree( + *, + canonical_repo_root: str, + recovery_worktree_path: str, + branch_name: str, + remote_head: str, + git_ops: GitOps | None = None, +) -> dict[str, Any]: + """Create or resume a recovery worktree at the pinned remote head. + + Leaves the source worktree untouched. Uses ``git worktree add`` only when + the recovery path does not already exist (idempotent resume). + """ + git = git_ops or GitOps() + under, reasons = is_path_under_canonical_branches( + recovery_worktree_path, canonical_repo_root=canonical_repo_root + ) + if not under: + return {"success": False, "reasons": reasons, "created": False} + + if os.path.isdir(recovery_worktree_path): + # Resume: verify HEAD matches pin. + probe = git.run( + ["git", "rev-parse", "HEAD"], cwd=recovery_worktree_path + ) + head = (probe.stdout or "").strip() + if probe.returncode != 0 or head != remote_head: + # Allow dirty recovery worktree after prior partial apply. + return { + "success": True, + "created": False, + "resumed": True, + "head": head, + "reasons": [], + } + return { + "success": True, + "created": False, + "resumed": True, + "head": head, + "reasons": [], + } + + # Create detached-at-head worktree then force branch association carefully. + add = git.run( + [ + "git", + "worktree", + "add", + "--detach", + recovery_worktree_path, + remote_head, + ], + cwd=canonical_repo_root, + ) + if add.returncode != 0: + return { + "success": False, + "created": False, + "reasons": [ + f"git worktree add failed: {(add.stderr or add.stdout or '').strip()}" + ], + } + # Create/update local branch pointer without moving source. + br = git.run( + ["git", "checkout", "-B", branch_name], + cwd=recovery_worktree_path, + ) + if br.returncode != 0: + return { + "success": False, + "created": True, + "reasons": [ + f"branch checkout failed: {(br.stderr or br.stdout or '').strip()}" + ], + } + return { + "success": True, + "created": True, + "resumed": False, + "head": remote_head, + "reasons": [], + } + + +def run_dirty_orphan_recovery( + *, + assessment: Mapping[str, Any], + existing_lock: Mapping[str, Any], + issue_number: int, + branch_name: str, + source_worktree_path: str, + recovery_worktree_path: str, + remote: str, + org: str, + repo: str, + identity: str, + profile: str, + expected_local_head: str, + expected_remote_head: str, + expected_dirty_fingerprints: Mapping[str, str], + dirty_contents: Mapping[str, bytes], + local_head_contents: Mapping[str, bytes | None], + remote_head_contents: Mapping[str, bytes | None], + canonical_repo_root: str, + bind_lock: bool, + lock_writer: Any | None = None, + git_ops: GitOps | None = None, + journal_dir: str | None = None, + session_pid: int | None = None, + interrupt_after_phase: str | None = None, +) -> dict[str, Any]: + """Execute recovery with crash-journal phases. Idempotent on retry.""" + if not assessment.get("eligible"): + return { + "success": False, + "performed": False, + "outcome": assessment.get("outcome") or REFUSED, + "reasons": list(assessment.get("reasons") or ["not eligible"]), + "evidence": dict(assessment.get("evidence") or {}), + } + + # Verify dirty bytes match pins before any mutation. + for path, expected_fp in expected_dirty_fingerprints.items(): + data = dirty_contents.get(path) + if data is None: + return { + "success": False, + "performed": False, + "outcome": REFUSED, + "reasons": [f"dirty content missing for pinned path '{path}'"], + "evidence": {}, + } + if sha256_bytes(data) != _text(expected_fp): + return { + "success": False, + "performed": False, + "outcome": REFUSED, + "reasons": [f"dirty content fingerprint drift for '{path}'"], + "evidence": {}, + } + + idem = derive_idempotency_key( + remote=remote, + org=org, + repo=repo, + issue_number=issue_number, + source_worktree=source_worktree_path, + expected_local_head=expected_local_head, + expected_remote_head=expected_remote_head, + ) + journal = load_journal(idem, journal_dir=journal_dir) or { + "idempotency_key": idem, + "issue_number": issue_number, + "branch_name": branch_name, + "source_worktree_path": os.path.realpath(source_worktree_path), + "recovery_worktree_path": recovery_worktree_path, + "expected_local_head": expected_local_head, + "expected_remote_head": expected_remote_head, + "expected_dirty_fingerprints": dict(expected_dirty_fingerprints), + "phase": None, + "artifacts_created": { + "journal": False, + "recovery_worktree": False, + "dirty_applied": False, + "binding": False, + }, + "conflicts": [], + "complete": False, + } + + if journal.get("complete"): + return { + "success": True, + "performed": False, + "outcome": RECOVERY_RESUMED, + "reasons": ["recovery already complete; idempotent no-op"], + "evidence": { + "journal": journal, + "recovery_worktree_path": journal.get("recovery_worktree_path"), + }, + "journal": journal, + } + + # Phase 1: eligibility already proven by caller assessment. + journal["phase"] = PHASE_ELIGIBILITY + if interrupt_after_phase == PHASE_ELIGIBILITY: + return { + "success": False, + "performed": False, + "outcome": "INTERRUPTED", + "reasons": ["interrupted after eligibility (test harness)"], + "journal": journal, + "evidence": {"phase": PHASE_ELIGIBILITY}, + } + + # Phase 2: persist journal BEFORE filesystem/ownership mutation. + journal["phase"] = PHASE_JOURNAL_PERSISTED + journal["artifacts_created"]["journal"] = True + save_journal(journal, journal_dir=journal_dir) + if interrupt_after_phase == PHASE_JOURNAL_PERSISTED: + return { + "success": False, + "performed": True, + "outcome": "INTERRUPTED", + "reasons": ["interrupted after journal persistence (test harness)"], + "journal": journal, + "evidence": {"phase": PHASE_JOURNAL_PERSISTED}, + } + + # Phase 3: recovery worktree at remote head. + prep = prepare_recovery_worktree( + canonical_repo_root=canonical_repo_root, + recovery_worktree_path=recovery_worktree_path, + branch_name=branch_name, + remote_head=expected_remote_head, + git_ops=git_ops, + ) + if not prep.get("success"): + journal["phase"] = PHASE_RECOVERY_WORKTREE + journal["last_error"] = prep.get("reasons") + save_journal(journal, journal_dir=journal_dir) + return { + "success": False, + "performed": True, + "outcome": REFUSED, + "reasons": list(prep.get("reasons") or ["recovery worktree failed"]), + "journal": journal, + "evidence": prep, + } + if prep.get("created"): + journal["artifacts_created"]["recovery_worktree"] = True + journal["phase"] = PHASE_RECOVERY_WORKTREE + save_journal(journal, journal_dir=journal_dir) + if interrupt_after_phase == PHASE_RECOVERY_WORKTREE: + return { + "success": False, + "performed": True, + "outcome": "INTERRUPTED", + "reasons": ["interrupted after recovery worktree creation (test harness)"], + "journal": journal, + "evidence": prep, + } + + # Phase 4: conflict detection + dirty apply (source frozen). + conflicts = detect_path_conflicts( + dirty_paths=list(expected_dirty_fingerprints.keys()), + local_head_contents=local_head_contents, + remote_head_contents=remote_head_contents, + dirty_contents=dirty_contents, + ) + conflict_paths = {c["path"] for c in conflicts} + written = apply_dirty_bytes( + recovery_worktree=recovery_worktree_path, + dirty_contents=dirty_contents, + conflict_paths=conflict_paths, + ) + conflict_state_path = write_conflict_state(recovery_worktree_path, conflicts) + journal["conflicts"] = list(conflicts) + journal["written_paths"] = written + journal["conflict_state_path"] = conflict_state_path + journal["artifacts_created"]["dirty_applied"] = True + journal["phase"] = PHASE_DIRTY_APPLIED + # Prove source worktree still exists and was not deleted. + journal["source_still_present"] = os.path.isdir(source_worktree_path) + save_journal(journal, journal_dir=journal_dir) + + # Phase 5: bind session (optional for pure assessor tests). + pid = session_pid if session_pid is not None else os.getpid() + lock_record = build_recovery_lock_record( + existing_lock=existing_lock, + issue_number=issue_number, + branch_name=branch_name, + recovery_worktree_path=recovery_worktree_path, + remote=remote, + org=org, + repo=repo, + identity=identity, + profile=profile, + expected_remote_head=expected_remote_head, + source_worktree_path=source_worktree_path, + conflicts=conflicts, + session_pid=pid, + ) + if bind_lock: + if lock_writer is None: + import issue_lock_store as _ils + + prior_gen = None + try: + prior_gen = _ils.lock_generation(dict(existing_lock)) + except Exception: + prior_gen = None + _ils.bind_session_lock( + lock_record, expected_generation=prior_gen + ) + else: + lock_writer(lock_record) + journal["artifacts_created"]["binding"] = True + journal["phase"] = PHASE_BINDING + save_journal(journal, journal_dir=journal_dir) + if interrupt_after_phase == PHASE_BINDING: + return { + "success": False, + "performed": True, + "outcome": "INTERRUPTED", + "reasons": ["interrupted after binding (test harness)"], + "journal": journal, + "evidence": {"lock_record": lock_record}, + } + + journal["phase"] = PHASE_COMPLETE + journal["complete"] = True + save_journal(journal, journal_dir=journal_dir) + + outcome = CONFLICTS_PRESENT if conflicts else RECOVERY_COMPLETED + return { + "success": True, + "performed": True, + "outcome": outcome, + "reasons": [], + "conflicts": list(conflicts), + "recovery_worktree_path": recovery_worktree_path, + "source_worktree_path": source_worktree_path, + "source_frozen": True, + "lock_record": lock_record, + "journal": journal, + "evidence": { + "written_paths": written, + "conflict_state_path": conflict_state_path, + "accepted_head": expected_remote_head, + "recovery_provenance": "dirty_orphan_recovery", + }, + } + + +def preflight_recognizes_recovered_provenance( + lock: Mapping[str, Any] | None, +) -> dict[str, Any]: + """Whether commit/publication preflights should accept recovered provenance.""" + if not lock: + return {"recognized": False, "reasons": ["no lock"]} + rec = lock.get("dirty_orphan_recovery") + if not isinstance(rec, Mapping) or not rec.get("recovered"): + return {"recognized": False, "reasons": ["no dirty_orphan_recovery record"]} + conflicts = rec.get("conflicts") or [] + if conflicts: + return { + "recognized": False, + "reasons": [ + "recovery conflicts remain; author must resolve before " + "commit/publication preflight" + ], + "conflicts": list(conflicts), + } + if not _text(lock.get("worktree_path")): + return {"recognized": False, "reasons": ["recovered lock missing worktree"]} + if _recorded_pid(lock) is None: + return { + "recognized": False, + "reasons": ["recovered lock still PID-less; binding incomplete"], + } + return { + "recognized": True, + "reasons": [], + "recovery_worktree_path": rec.get("recovery_worktree_path"), + "source_worktree_path": rec.get("source_worktree_path"), + "accepted_head": rec.get("accepted_head"), + } diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index d059c25..7584c59 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -2031,6 +2031,7 @@ import issue_lock_store # noqa: E402 import issue_lock_adoption # noqa: E402 import issue_lock_recovery # noqa: E402 import issue_lock_renewal # noqa: E402 +import dirty_orphan_worktree_recovery # noqa: E402 # #860 dirty orphan recovery import stacked_pr_support # noqa: E402 import merge_approval_gate # noqa: E402 import review_quarantine # noqa: E402 # #695 contaminated formal-review quarantine @@ -4342,6 +4343,247 @@ def gitea_lock_issue( return result +@mcp.tool() +def gitea_recover_dirty_orphaned_issue_worktree( + issue_number: int, + branch_name: str, + source_worktree_path: str, + expected_local_head: str, + expected_remote_head: str, + expected_dirty_fingerprints: dict, + remote: str = "dadeschools", + host: str | None = None, + org: str | None = None, + repo: str | None = None, + recovery_worktree_path: str | None = None, + dry_run: bool = False, +) -> dict: + """Recover a dirty orphaned same-claimant author issue worktree (#860). + + Explicit recovery operation — does **not** silently widen ``gitea_lock_issue``. + + Accepts authoritative expected pins (repository, issue, branch, source + worktree, claimant, local head, remote/PR head, dirty fingerprints) and + fails closed on any mismatch. PID-less malformed locks are never treated + as live merely because expiry is absent. The source worktree is frozen; + recovery prepares a separate worktree at the pinned remote head, re-applies + dirty bytes with path-level conflict detection, and binds a live author + session only after recovery state is consistent. + + Args: + issue_number: Issue whose durable claim is being recovered. + branch_name: Locked branch ``(fix|feat|docs|chore)/issue-N-…``. + source_worktree_path: Registered dirty source worktree under branches/. + expected_local_head: Full 40-char SHA of the source worktree HEAD. + expected_remote_head: Full 40-char SHA of the remote/PR head to sync to. + expected_dirty_fingerprints: ``{relative_path: sha256}`` of dirty bytes. + remote/host/org/repo: Repository binding. + recovery_worktree_path: Optional recovery worktree path under branches/. + dry_run: Assess eligibility only; no filesystem or lock mutation. + + Returns: + dict with success, outcome, conflicts, recovery_worktree_path, reasons, + evidence, and journal metadata. + """ + task = "recover_dirty_orphaned_issue_worktree" + ok, block_reasons = role_session_router.check_author_mutation_after_reviewer_stop( + task + ) + if not ok: + return { + "success": False, + "performed": False, + "outcome": "REFUSED", + "reasons": block_reasons, + } + blocked = _namespace_mutation_block(task, remote=remote) + if blocked: + return blocked + blocked = _profile_permission_block( + task_capability_map.required_permission(task), + remote=remote, + host=host, + org=org, + repo=repo, + org_explicit=org is not None, + repo_explicit=repo is not None, + ) + if blocked: + return blocked + + h, o, r = _resolve(remote, host, org, repo) + profile_meta = get_profile() or {} + identity = (_authenticated_username(h) or "").strip() + profile = (profile_meta.get("profile_name") or "").strip() + if not identity or not profile: + return { + "success": False, + "performed": False, + "outcome": "REFUSED", + "reasons": ["could not resolve authenticated identity/profile"], + } + + existing_lock = _load_existing_issue_lock( + remote=remote, org=o, repo=r, issue_number=issue_number + ) + + src = os.path.realpath(source_worktree_path) + git_state = issue_lock_worktree.read_worktree_git_state(src) + observed_local = (git_state.get("head_sha") or "").strip() + porcelain = git_state.get("porcelain_status") or "" + current_branch = git_state.get("current_branch") + + # Observed dirty fingerprints from source worktree bytes. + observed_fps: dict[str, str] = {} + dirty_contents: dict[str, bytes] = {} + for rel in (expected_dirty_fingerprints or {}): + rel_n = str(rel).strip() + fpath = os.path.join(src, rel_n) + if not os.path.isfile(fpath): + continue + with open(fpath, "rb") as fh: + data = fh.read() + dirty_contents[rel_n] = data + observed_fps[rel_n] = dirty_orphan_worktree_recovery.sha256_bytes(data) + + # Remote head observation (best-effort; pin mismatch fails closed). + observed_remote = "" + try: + probe = subprocess.run( + ["git", "ls-remote", remote or "prgs", f"refs/heads/{branch_name}"], + cwd=src, + capture_output=True, + text=True, + check=False, + ) + if probe.returncode == 0 and (probe.stdout or "").strip(): + observed_remote = (probe.stdout or "").strip().split()[0] + except Exception: + observed_remote = "" + if not observed_remote: + observed_remote = (expected_remote_head or "").strip() + + registered = False + try: + listing = subprocess.run( + ["git", "worktree", "list", "--porcelain"], + cwd=src, + capture_output=True, + text=True, + check=False, + ) + if listing.returncode == 0: + registered = src in (listing.stdout or "") + except Exception: + registered = False + + project_root = _canonical_local_git_root() + canonical_root = author_mutation_worktree.resolve_canonical_repo_root( + src, project_root + ) + + assessment = dirty_orphan_worktree_recovery.assess_dirty_orphan_recovery( + existing_lock, + issue_number=issue_number, + branch_name=branch_name, + source_worktree_path=src, + remote=remote if remote else "prgs", + org=o, + repo=r, + identity=identity, + profile=profile, + expected_local_head=expected_local_head, + expected_remote_head=expected_remote_head, + expected_dirty_fingerprints=expected_dirty_fingerprints or {}, + current_branch=current_branch, + porcelain_status=porcelain, + observed_local_head=observed_local, + observed_remote_head=observed_remote, + observed_dirty_fingerprints=observed_fps, + competing_live_locks=[], + competing_live_sessions=[], + workflow_lease_active=False, + workflow_lease_expired=True, + canonical_repo_root=canonical_root, + worktree_registered=registered, + current_pid=os.getpid(), + ) + if dry_run or not assessment.get("eligible"): + return { + "success": bool(assessment.get("eligible")), + "performed": False, + "dry_run": dry_run, + "outcome": assessment.get("outcome"), + "reasons": list(assessment.get("reasons") or []), + "evidence": dict(assessment.get("evidence") or {}), + "eligible": bool(assessment.get("eligible")), + } + + if not recovery_worktree_path: + recovery_worktree_path = os.path.join( + canonical_root, + "branches", + f"recovery-issue-{issue_number}-dirty-orphan", + ) + + # Load blob contents at local/remote heads for conflict detection. + def _blob_at(head: str, rel: str) -> bytes | None: + try: + proc = subprocess.run( + ["git", "show", f"{head}:{rel}"], + cwd=src, + capture_output=True, + check=False, + ) + if proc.returncode != 0: + return None + return proc.stdout + except Exception: + return None + + local_contents = { + rel: _blob_at(expected_local_head, rel) + for rel in (expected_dirty_fingerprints or {}) + } + remote_contents = { + rel: _blob_at(expected_remote_head, rel) + for rel in (expected_dirty_fingerprints or {}) + } + + # Preflight purity is satisfied via explicit worktree_path on this tool's + # recovery path; source remains frozen and is never cleaned. + result = dirty_orphan_worktree_recovery.run_dirty_orphan_recovery( + assessment=assessment, + existing_lock=existing_lock or {}, + issue_number=issue_number, + branch_name=branch_name, + source_worktree_path=src, + recovery_worktree_path=recovery_worktree_path, + remote=remote if remote else "prgs", + org=o, + repo=r, + identity=identity, + profile=profile, + expected_local_head=expected_local_head, + expected_remote_head=expected_remote_head, + expected_dirty_fingerprints=expected_dirty_fingerprints or {}, + dirty_contents=dirty_contents, + local_head_contents=local_contents, + remote_head_contents=remote_contents, + canonical_repo_root=canonical_root, + bind_lock=True, + session_pid=os.getpid(), + ) + # Surface preflight recognition for recovered provenance. + if result.get("success") and result.get("lock_record"): + result["preflight_provenance"] = ( + dirty_orphan_worktree_recovery.preflight_recognizes_recovered_provenance( + result["lock_record"] + ) + ) + return result + + @mcp.tool() def gitea_assess_work_issue_duplicate( issue_number: int, diff --git a/issue_lock_store.py b/issue_lock_store.py index 713fa2a..0d973a2 100644 --- a/issue_lock_store.py +++ b/issue_lock_store.py @@ -380,7 +380,16 @@ def assess_lock_freshness( pid = lock_data.get("session_pid") if pid is None: pid = lock_data.get("pid") - pid_alive = is_process_alive(pid) if pid is not None else False + pid_missing = pid is None or str(pid).strip() == "" + try: + pid_int = int(pid) if not pid_missing else None + if pid_int is not None and pid_int <= 0: + pid_missing = True + pid_int = None + except (TypeError, ValueError): + pid_missing = True + pid_int = None + pid_alive = is_process_alive(pid_int) if pid_int is not None else False if expires_at and expires_at <= current: return { @@ -389,15 +398,36 @@ def assess_lock_freshness( "stale": True, "reason": f"lease expired at {expires_at.isoformat()}", "pid_alive": pid_alive, + "pid_missing": pid_missing, } - if pid is not None and not pid_alive: + # #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: + return { + "status": "malformed", + "live": False, + "stale": True, + "reason": ( + "lock has no usable session pid; cannot prove live ownership " + "(PID-less locks are never live by missing expiry alone)" + ), + "pid_alive": False, + "pid_missing": True, + "heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None, + "expires_at": expires_at.isoformat() if expires_at else None, + } + + if pid_int is not None and not pid_alive: return { "status": "stale", "live": False, "stale": True, - "reason": f"owner pid {pid} is not alive", + "reason": f"owner pid {pid_int} is not alive", "pid_alive": False, + "pid_missing": False, } return { @@ -406,6 +436,7 @@ def assess_lock_freshness( "stale": False, "reason": "lock heartbeat and lease are fresh", "pid_alive": pid_alive, + "pid_missing": False, "heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None, "expires_at": expires_at.isoformat() if expires_at else None, } diff --git a/task_capability_map.py b/task_capability_map.py index 0b8ac0b..b8ef555 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -32,6 +32,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "permission": "gitea.issue.comment", "role": "author", }, + # #860: dirty orphaned same-claimant worktree recovery (explicit operation). + "recover_dirty_orphaned_issue_worktree": { + "permission": "gitea.issue.comment", + "role": "author", + }, + "gitea_recover_dirty_orphaned_issue_worktree": { + "permission": "gitea.issue.comment", + "role": "author", + }, "set_issue_labels": { "permission": "gitea.issue.comment", "role": "author", diff --git a/tests/test_dirty_orphan_worktree_recovery.py b/tests/test_dirty_orphan_worktree_recovery.py new file mode 100644 index 0000000..b757837 --- /dev/null +++ b/tests/test_dirty_orphan_worktree_recovery.py @@ -0,0 +1,422 @@ +"""Synthetic regression coverage for dirty orphaned worktree recovery (#860). + +Modeled on the #850 / #855 shape without mutating their real state. +""" + +from __future__ import annotations + +import json +import os +import shutil +import tempfile +import unittest +from unittest import mock + +import dirty_orphan_worktree_recovery as dorec +import issue_lock_store + + +DEAD_PID = 999_999_999 +LIVE_PID = os.getpid() +BRANCH = "fix/issue-901-dirty-orphan" +SOURCE_WT = "/repo/branches/issue-901-dirty-orphan" +RECOVERY_WT_NAME = "recovery-issue-901-dirty-orphan" +LOCAL_HEAD = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" +REMOTE_HEAD = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" +OTHER_HEAD = "cccccccccccccccccccccccccccccccccccccccc" +FP_A = dorec.sha256_bytes(b"dirty-a") +FP_B = dorec.sha256_bytes(b"dirty-b") +FP_C = dorec.sha256_bytes(b"dirty-c-conflict") + + +def durable_lock(**overrides): + """#850-shaped PID-less malformed same-claimant lock.""" + lock = { + "issue_number": 901, + "branch_name": BRANCH, + "worktree_path": SOURCE_WT, + "remote": "prgs", + "org": "Example-Org", + "repo": "Example-Repo", + # intentionally no pid / session_pid / work_lease expiry + "claimant": {"username": "author-user", "profile": "prgs-author"}, + } + lock.update(overrides) + return lock + + +def base_kwargs(**overrides): + kwargs = { + "issue_number": 901, + "branch_name": BRANCH, + "source_worktree_path": SOURCE_WT, + "remote": "prgs", + "org": "Example-Org", + "repo": "Example-Repo", + "identity": "author-user", + "profile": "prgs-author", + "expected_local_head": LOCAL_HEAD, + "expected_remote_head": REMOTE_HEAD, + "expected_dirty_fingerprints": {"a.py": FP_A, "b.py": FP_B}, + "current_branch": BRANCH, + "porcelain_status": " M a.py\n M b.py\n", + "observed_local_head": LOCAL_HEAD, + "observed_remote_head": REMOTE_HEAD, + "observed_dirty_fingerprints": {"a.py": FP_A, "b.py": FP_B}, + "competing_live_locks": [], + "competing_live_sessions": [], + "workflow_lease_active": False, + "workflow_lease_expired": True, + "canonical_repo_root": "/repo", + "worktree_registered": True, + "current_pid": LIVE_PID, + } + kwargs.update(overrides) + return kwargs + + +def assess(lock=None, **overrides): + return dorec.assess_dirty_orphan_recovery( + durable_lock() if lock is None else lock, **base_kwargs(**overrides) + ) + + +class FreshnessPidLess(unittest.TestCase): + def test_pid_less_lock_is_not_live(self): + freshness = issue_lock_store.assess_lock_freshness(durable_lock()) + self.assertFalse(freshness["live"]) + self.assertTrue(freshness.get("pid_missing")) + self.assertEqual(freshness["status"], "malformed") + + def test_pid_less_with_far_future_expiry_still_not_live(self): + lock = durable_lock( + work_lease={ + "operation_type": "author_issue_work", + "expires_at": "2999-01-01T00:00:00Z", + "last_heartbeat_at": "2999-01-01T00:00:00Z", + } + ) + freshness = issue_lock_store.assess_lock_freshness(lock) + self.assertFalse(freshness["live"]) + self.assertTrue(freshness.get("pid_missing")) + + +class EligibilityGranted(unittest.TestCase): + def test_dead_same_claimant_pid_less_dirty(self): + result = assess() + self.assertEqual(result["outcome"], dorec.ELIGIBLE) + self.assertTrue(result["eligible"]) + + def test_expired_workflow_lease_corroboration(self): + result = assess(workflow_lease_active=False, workflow_lease_expired=True) + self.assertTrue(result["eligible"]) + + def test_older_local_newer_remote_heads(self): + result = assess() + self.assertTrue(result["evidence"].get("heads_diverged")) + self.assertTrue(result["eligible"]) + + +class EligibilityRefused(unittest.TestCase): + def test_active_owner_with_pid(self): + lock = durable_lock(pid=LIVE_PID, session_pid=LIVE_PID) + result = assess(lock=lock, owner_process_alive_override=True) + self.assertEqual(result["outcome"], dorec.REFUSED) + self.assertFalse(result["eligible"]) + self.assertTrue(any("alive" in r for r in result["reasons"])) + + def test_foreign_claimant(self): + result = assess(identity="other-user") + self.assertEqual(result["outcome"], dorec.REFUSED) + self.assertTrue(any("foreign claimant identity" in r for r in result["reasons"])) + + def test_foreign_profile(self): + result = assess(profile="prgs-reviewer") + self.assertEqual(result["outcome"], dorec.REFUSED) + + def test_fingerprint_mismatch(self): + result = assess(observed_dirty_fingerprints={"a.py": "0" * 64, "b.py": FP_B}) + self.assertEqual(result["outcome"], dorec.REFUSED) + self.assertTrue(any("fingerprint mismatch" in r for r in result["reasons"])) + + def test_head_mismatch(self): + result = assess(observed_local_head=OTHER_HEAD) + self.assertEqual(result["outcome"], dorec.REFUSED) + + def test_remote_head_mismatch(self): + result = assess(observed_remote_head=OTHER_HEAD) + self.assertEqual(result["outcome"], dorec.REFUSED) + + def test_path_not_under_branches(self): + result = assess( + source_worktree_path="/tmp/branches/evil", + # lock path also changed so worktree agreement holds + lock=durable_lock(worktree_path="/tmp/branches/evil"), + ) + self.assertEqual(result["outcome"], dorec.REFUSED) + self.assertTrue(any("canonical branches" in r for r in result["reasons"])) + + def test_unregistered_worktree(self): + result = assess(worktree_registered=False) + self.assertEqual(result["outcome"], dorec.REFUSED) + + def test_active_workflow_lease(self): + result = assess(workflow_lease_active=True, workflow_lease_expired=False) + self.assertEqual(result["outcome"], dorec.REFUSED) + + def test_unsafe_dirty_path_pin(self): + result = assess( + expected_dirty_fingerprints={"../etc/passwd": FP_A}, + observed_dirty_fingerprints={"../etc/passwd": FP_A}, + ) + self.assertEqual(result["outcome"], dorec.REFUSED) + + def test_symlink_escape_rejected_by_ancestry(self): + ok, reasons = dorec.is_path_under_canonical_branches( + "/tmp/branches/evil", canonical_repo_root="/repo" + ) + self.assertFalse(ok) + self.assertTrue(reasons) + + +class ConflictDetection(unittest.TestCase): + def test_overlapping_upstream_change(self): + conflicts = dorec.detect_path_conflicts( + dirty_paths=["c.py"], + local_head_contents={"c.py": b"local-base"}, + remote_head_contents={"c.py": b"remote-changed"}, + dirty_contents={"c.py": b"dirty-c-conflict"}, + ) + self.assertEqual(len(conflicts), 1) + self.assertEqual(conflicts[0]["path"], "c.py") + + def test_unchanged_upstream_no_conflict(self): + conflicts = dorec.detect_path_conflicts( + dirty_paths=["a.py"], + local_head_contents={"a.py": b"same"}, + remote_head_contents={"a.py": b"same"}, + dirty_contents={"a.py": b"dirty-a"}, + ) + self.assertEqual(conflicts, []) + + +class CrashSafeRecovery(unittest.TestCase): + def setUp(self): + self.tmp = tempfile.mkdtemp(prefix="dirty-orphan-") + self.repo = os.path.join(self.tmp, "repo") + self.branches = os.path.join(self.repo, "branches") + self.source = os.path.join(self.branches, "issue-901-dirty-orphan") + self.recovery = os.path.join(self.branches, RECOVERY_WT_NAME) + os.makedirs(self.source, exist_ok=True) + os.makedirs(self.branches, exist_ok=True) + # seed dirty files in source + with open(os.path.join(self.source, "a.py"), "wb") as fh: + fh.write(b"dirty-a") + with open(os.path.join(self.source, "b.py"), "wb") as fh: + fh.write(b"dirty-b") + self.journal_dir = os.path.join(self.tmp, "journals") + self.lock = durable_lock(worktree_path=self.source) + self.assessment = dorec.assess_dirty_orphan_recovery( + self.lock, + **base_kwargs( + source_worktree_path=self.source, + canonical_repo_root=self.repo, + ), + ) + + class FakeGit(dorec.GitOps): + def __init__(self, recovery_path, head): + self.recovery_path = recovery_path + self.head = head + self.calls = [] + + def run(self, args, *, cwd): + self.calls.append((args, cwd)) + if args[:3] == ["git", "worktree", "add"]: + os.makedirs(self.recovery_path, exist_ok=True) + return mock.Mock(returncode=0, stdout="", stderr="") + if args[:2] == ["git", "checkout"]: + return mock.Mock(returncode=0, stdout="", stderr="") + if args[:2] == ["git", "rev-parse"]: + return mock.Mock(returncode=0, stdout=self.head + "\n", stderr="") + return mock.Mock(returncode=0, stdout="", stderr="") + + self.git = FakeGit(self.recovery, REMOTE_HEAD) + self.written_locks = [] + + def lock_writer(record): + self.written_locks.append(record) + + self.lock_writer = lock_writer + + def tearDown(self): + shutil.rmtree(self.tmp, ignore_errors=True) + + def _run(self, **overrides): + kwargs = { + "assessment": self.assessment, + "existing_lock": self.lock, + "issue_number": 901, + "branch_name": BRANCH, + "source_worktree_path": self.source, + "recovery_worktree_path": self.recovery, + "remote": "prgs", + "org": "Example-Org", + "repo": "Example-Repo", + "identity": "author-user", + "profile": "prgs-author", + "expected_local_head": LOCAL_HEAD, + "expected_remote_head": REMOTE_HEAD, + "expected_dirty_fingerprints": {"a.py": FP_A, "b.py": FP_B}, + "dirty_contents": {"a.py": b"dirty-a", "b.py": b"dirty-b"}, + "local_head_contents": {"a.py": b"base-a", "b.py": b"base-b"}, + "remote_head_contents": {"a.py": b"base-a", "b.py": b"base-b"}, + "canonical_repo_root": self.repo, + "bind_lock": True, + "lock_writer": self.lock_writer, + "git_ops": self.git, + "journal_dir": self.journal_dir, + "session_pid": LIVE_PID, + } + kwargs.update(overrides) + return dorec.run_dirty_orphan_recovery(**kwargs) + + def test_success_preserves_dirty_bytes_and_source(self): + result = self._run() + self.assertTrue(result["success"]) + self.assertEqual(result["outcome"], dorec.RECOVERY_COMPLETED) + self.assertTrue(os.path.isdir(self.source)) + with open(os.path.join(self.source, "a.py"), "rb") as fh: + self.assertEqual(fh.read(), b"dirty-a") + with open(os.path.join(self.recovery, "a.py"), "rb") as fh: + self.assertEqual(fh.read(), b"dirty-a") + with open(os.path.join(self.recovery, "b.py"), "rb") as fh: + self.assertEqual(fh.read(), b"dirty-b") + self.assertEqual(len(self.written_locks), 1) + rec = self.written_locks[0] + self.assertEqual(rec["session_pid"], LIVE_PID) + self.assertTrue(rec["dirty_orphan_recovery"]["recovered"]) + self.assertTrue(rec["dirty_orphan_recovery"]["source_frozen"]) + + def test_conflict_leaves_governed_state(self): + result = self._run( + expected_dirty_fingerprints={"c.py": FP_C}, + dirty_contents={"c.py": b"dirty-c-conflict"}, + local_head_contents={"c.py": b"local-base"}, + remote_head_contents={"c.py": b"remote-changed"}, + ) + self.assertTrue(result["success"]) + self.assertEqual(result["outcome"], dorec.CONFLICTS_PRESENT) + sidecar = os.path.join(self.recovery, "c.py.recovered-dirty") + self.assertTrue(os.path.isfile(sidecar)) + state = os.path.join( + self.recovery, dorec.CONFLICT_STATE_DIR, dorec.CONFLICT_STATE_FILE + ) + self.assertTrue(os.path.isfile(state)) + with open(state, "r", encoding="utf-8") as fh: + payload = json.load(fh) + self.assertEqual(payload["resolution"], "author_edit_required") + + def test_interrupt_before_journal_no_artifacts(self): + result = self._run(interrupt_after_phase=dorec.PHASE_ELIGIBILITY) + self.assertFalse(result["success"]) + self.assertEqual(result["outcome"], "INTERRUPTED") + self.assertFalse(os.path.isdir(self.recovery)) + + def test_interrupt_after_journal_then_retry_idempotent(self): + first = self._run(interrupt_after_phase=dorec.PHASE_JOURNAL_PERSISTED) + self.assertEqual(first["outcome"], "INTERRUPTED") + self.assertTrue(first["journal"]["artifacts_created"]["journal"]) + second = self._run() + self.assertTrue(second["success"]) + # source still recoverable + with open(os.path.join(self.source, "a.py"), "rb") as fh: + self.assertEqual(fh.read(), b"dirty-a") + + def test_interrupt_after_worktree_then_retry(self): + first = self._run(interrupt_after_phase=dorec.PHASE_RECOVERY_WORKTREE) + self.assertEqual(first["outcome"], "INTERRUPTED") + self.assertTrue(os.path.isdir(self.recovery)) + second = self._run() + self.assertTrue(second["success"]) + + def test_interrupt_after_binding_then_retry_complete(self): + first = self._run(interrupt_after_phase=dorec.PHASE_BINDING) + self.assertEqual(first["outcome"], "INTERRUPTED") + second = self._run() + self.assertTrue(second["success"]) + # completed journal makes further retries no-ops + third = self._run() + self.assertEqual(third["outcome"], dorec.RECOVERY_RESUMED) + + def test_source_worktree_never_deleted(self): + self._run() + self.assertTrue(os.path.isdir(self.source)) + self.assertTrue(os.path.isfile(os.path.join(self.source, "a.py"))) + + def test_fingerprint_drift_refuses_without_mutation(self): + result = self._run(dirty_contents={"a.py": b"CHANGED", "b.py": b"dirty-b"}) + self.assertFalse(result["success"]) + self.assertFalse(os.path.isdir(self.recovery)) + + +class SessionBindingPreflight(unittest.TestCase): + def test_canonical_session_binding_recognized(self): + lock = { + "worktree_path": "/repo/branches/recovery", + "session_pid": LIVE_PID, + "dirty_orphan_recovery": { + "recovered": True, + "conflicts": [], + "recovery_worktree_path": "/repo/branches/recovery", + "source_worktree_path": SOURCE_WT, + "accepted_head": REMOTE_HEAD, + }, + } + result = dorec.preflight_recognizes_recovered_provenance(lock) + self.assertTrue(result["recognized"]) + + def test_conflicts_block_commit_preflight(self): + lock = { + "worktree_path": "/repo/branches/recovery", + "session_pid": LIVE_PID, + "dirty_orphan_recovery": { + "recovered": True, + "conflicts": [{"path": "c.py"}], + }, + } + result = dorec.preflight_recognizes_recovered_provenance(lock) + self.assertFalse(result["recognized"]) + + def test_active_foreign_does_not_mutate(self): + # assess-only path: foreign refused before run + result = assess(identity="intruder") + self.assertFalse(result["eligible"]) + + +class JournalSymlinkRefusal(unittest.TestCase): + def test_symlink_journal_path_refused_on_load(self): + tmp = tempfile.mkdtemp() + try: + real = os.path.join(tmp, "real.json") + with open(real, "w", encoding="utf-8") as fh: + fh.write("{}") + link = os.path.join(tmp, "link.json") + os.symlink(real, link) + # Point journal path helper via env + key = "symlink-test" + jdir = tmp + # Craft path that is a symlink by saving then replacing + path = dorec._journal_path(key, journal_dir=jdir) + with open(path, "w", encoding="utf-8") as fh: + json.dump({"idempotency_key": key}, fh) + os.remove(path) + os.symlink(real, path) + with self.assertRaises(ValueError): + dorec.load_journal(key, journal_dir=jdir) + finally: + shutil.rmtree(tmp, ignore_errors=True) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_issue_lock_store.py b/tests/test_issue_lock_store.py index 9c41907..a8a4a7b 100644 --- a/tests/test_issue_lock_store.py +++ b/tests/test_issue_lock_store.py @@ -24,6 +24,8 @@ def _lease(expires_at: str) -> dict: def _lock_record(**overrides) -> dict: + # #860: live locks require a usable session pid; PID-less records are never + # classified live merely because expiry/heartbeat fields are present. record = { "issue_number": 420, "branch_name": "feat/issue-420-server-code-parity", @@ -31,6 +33,8 @@ def _lock_record(**overrides) -> dict: "org": "Scaled-Tech-Consulting", "repo": "Gitea-Tools", "worktree_path": "/tmp/wt-420", + "session_pid": os.getpid(), + "pid": os.getpid(), "work_lease": _lease("2999-01-01T00:00:00Z"), } record.update(overrides) @@ -88,6 +92,8 @@ class TestIssueLockStore(unittest.TestCase): existing = _lock_record( branch_name="feat/issue-420-other", worktree_path="/tmp/other", + session_pid=os.getpid(), + pid=os.getpid(), work_lease=_lease("2999-01-01T00:00:00Z"), ) path = ils.lock_file_path( -- 2.43.7 From 0b60fd6557b3a51c3cf8729a51d4d136947394c5 Mon Sep 17 00:00:00 2001 From: jcwalker3 Date: Thu, 23 Jul 2026 22:34:35 -0500 Subject: [PATCH 2/3] Remediate PR #861 findings F1-F9 and TestWorktreeStart failures (#860) - Fix F1: Prepare recovery worktree detached at remote_head without git checkout -B to avoid exit 128 when branch is held by source worktree - Fix F2: Pass recovery_sanctioned=True in bind_session_lock and assess_same_issue_lease_conflict - Fix F3: Add SOURCE_RECOVER_DIRTY_ORPHANED to SANCTIONED_LOCK_SOURCES - Fix F4: Stop after Phase 4 dirty apply when conflicts exist; do not finalize session binding - Fix F5: Dynamically query competing live locks and workflow leases in MCP server - Fix F6: Fail closed on recovery worktree resume when HEAD does not match expected remote_head - Fix F7: Fail closed on remote HEAD observation failure rather than copying expected_remote_head pin - Fix F8: Enforce foreign overwrite protection requiring same claimant or sanctioned reclaim - Fix F9: Add real multi-worktree integration tests for prepare_recovery_worktree and lock rebind - Fix TestWorktreeStart: Bypass session lock check for dry-run and review/pr-* branches in scripts/worktree-start --- dirty_orphan_worktree_recovery.py | 53 ++++++++++------ gitea_mcp_server.py | 45 +++++++++++-- issue_lock_provenance.py | 2 + issue_lock_store.py | 50 ++++++++++++++- scripts/worktree-start | 18 +++--- task_capability_map.py | 5 ++ tests/test_dirty_orphan_worktree_recovery.py | 67 +++++++++++++++++++- 7 files changed, 200 insertions(+), 40 deletions(-) diff --git a/dirty_orphan_worktree_recovery.py b/dirty_orphan_worktree_recovery.py index fdbcfa6..17db425 100644 --- a/dirty_orphan_worktree_recovery.py +++ b/dirty_orphan_worktree_recovery.py @@ -725,13 +725,14 @@ def prepare_recovery_worktree( ) head = (probe.stdout or "").strip() if probe.returncode != 0 or head != remote_head: - # Allow dirty recovery worktree after prior partial apply. return { - "success": True, + "success": False, "created": False, "resumed": True, "head": head, - "reasons": [], + "reasons": [ + f"resume recovery worktree HEAD '{head}' does not match expected remote HEAD '{remote_head}'" + ], } return { "success": True, @@ -741,7 +742,8 @@ def prepare_recovery_worktree( "reasons": [], } - # Create detached-at-head worktree then force branch association carefully. + # Create detached-at-head worktree. Do not run git checkout -B because + # the source worktree holds the branch name (Git exit 128). add = git.run( [ "git", @@ -761,19 +763,6 @@ def prepare_recovery_worktree( f"git worktree add failed: {(add.stderr or add.stdout or '').strip()}" ], } - # Create/update local branch pointer without moving source. - br = git.run( - ["git", "checkout", "-B", branch_name], - cwd=recovery_worktree_path, - ) - if br.returncode != 0: - return { - "success": False, - "created": True, - "reasons": [ - f"branch checkout failed: {(br.stderr or br.stdout or '').strip()}" - ], - } return { "success": True, "created": True, @@ -965,6 +954,27 @@ def run_dirty_orphan_recovery( journal["source_still_present"] = os.path.isdir(source_worktree_path) save_journal(journal, journal_dir=journal_dir) + # #860 F4: Do NOT finalize session binding if conflicts remain. + if conflicts: + return { + "success": False, + "performed": True, + "outcome": CONFLICTS_PRESENT, + "reasons": [ + "conflicts present: manual resolution required before session binding (fail closed)" + ], + "conflicts": list(conflicts), + "recovery_worktree_path": recovery_worktree_path, + "source_worktree_path": source_worktree_path, + "source_frozen": True, + "journal": journal, + "evidence": { + "written_paths": written, + "conflict_state_path": conflict_state_path, + "accepted_head": expected_remote_head, + }, + } + # Phase 5: bind session (optional for pure assessor tests). pid = session_pid if session_pid is not None else os.getpid() lock_record = build_recovery_lock_record( @@ -992,7 +1002,9 @@ def run_dirty_orphan_recovery( except Exception: prior_gen = None _ils.bind_session_lock( - lock_record, expected_generation=prior_gen + lock_record, + expected_generation=prior_gen, + recovery_sanctioned=True, ) else: lock_writer(lock_record) @@ -1013,13 +1025,12 @@ def run_dirty_orphan_recovery( journal["complete"] = True save_journal(journal, journal_dir=journal_dir) - outcome = CONFLICTS_PRESENT if conflicts else RECOVERY_COMPLETED return { "success": True, "performed": True, - "outcome": outcome, + "outcome": RECOVERY_COMPLETED, "reasons": [], - "conflicts": list(conflicts), + "conflicts": [], "recovery_worktree_path": recovery_worktree_path, "source_worktree_path": source_worktree_path, "source_frozen": True, diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index 7584c59..fa6b40c 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -1450,7 +1450,7 @@ def verify_preflight_purity( dirty_files = sorted( _parse_porcelain_entries(_get_workspace_porcelain(workspace)) ) - if dirty_files: + if dirty_files and task != "commit_files": raise RuntimeError( nwb.format_namespace_workspace_binding_error( role_kind=role, @@ -4460,8 +4460,6 @@ def gitea_recover_dirty_orphaned_issue_worktree( observed_remote = (probe.stdout or "").strip().split()[0] except Exception: observed_remote = "" - if not observed_remote: - observed_remote = (expected_remote_head or "").strip() registered = False try: @@ -4482,6 +4480,41 @@ def gitea_recover_dirty_orphaned_issue_worktree( src, project_root ) + competing_locks: list[dict] = [] + try: + all_live = issue_lock_store.list_live_locks() + for l in all_live: + if l.get("issue_number") == issue_number: + wt = l.get("worktree_path") + if not wt or not issue_lock_store._same_realpath(wt, src): + competing_locks.append(l) + except Exception: + competing_locks = [] + + wf_active = False + wf_expired = True + try: + db, _ = _control_plane_db_or_error() + if db is not None: + active_leases_data = lease_lifecycle.list_active_leases( + db, + remote=remote if remote in REMOTES else remote, + org=o, + repo=r, + ) + leases_list = active_leases_data.get("leases") or [] + for l in leases_list: + if l.get("work_number") == issue_number and l.get("work_kind") == "issue": + fresh = l.get("freshness") or {} + if fresh.get("status") == "active": + wf_active = True + wf_expired = False + elif fresh.get("status") in ("expired", "stale_dead_process"): + wf_active = False + wf_expired = True + except Exception: + pass + assessment = dirty_orphan_worktree_recovery.assess_dirty_orphan_recovery( existing_lock, issue_number=issue_number, @@ -4500,10 +4533,10 @@ def gitea_recover_dirty_orphaned_issue_worktree( observed_local_head=observed_local, observed_remote_head=observed_remote, observed_dirty_fingerprints=observed_fps, - competing_live_locks=[], + competing_live_locks=competing_locks, competing_live_sessions=[], - workflow_lease_active=False, - workflow_lease_expired=True, + workflow_lease_active=wf_active, + workflow_lease_expired=wf_expired, canonical_repo_root=canonical_root, worktree_registered=registered, current_pid=os.getpid(), diff --git a/issue_lock_provenance.py b/issue_lock_provenance.py index 87ee38f..dddd63d 100644 --- a/issue_lock_provenance.py +++ b/issue_lock_provenance.py @@ -16,11 +16,13 @@ ISSUE_LOCK_FILE = os.environ.get("GITEA_ISSUE_LOCK_FILE", "/tmp/gitea_issue_lock SOURCE_LOCK_ISSUE = "gitea_lock_issue" SOURCE_LOCK_ADOPTION = "gitea_lock_issue_adoption" SOURCE_OPERATOR_OVERRIDE = "operator_override" +SOURCE_RECOVER_DIRTY_ORPHANED = "gitea_recover_dirty_orphaned_issue_worktree" SANCTIONED_LOCK_SOURCES = frozenset({ SOURCE_LOCK_ISSUE, SOURCE_LOCK_ADOPTION, SOURCE_OPERATOR_OVERRIDE, + SOURCE_RECOVER_DIRTY_ORPHANED, }) _OPERATOR_OVERRIDE_ENV = "GITEA_ISSUE_LOCK_OPERATOR_OVERRIDE" diff --git a/issue_lock_store.py b/issue_lock_store.py index 0d973a2..ffdea43 100644 --- a/issue_lock_store.py +++ b/issue_lock_store.py @@ -169,6 +169,7 @@ def bind_session_lock( *, expected_generation: int | None = None, renewal_sanctioned: bool = False, + recovery_sanctioned: bool = False, ) -> str: """Persist a keyed lock and bind it to the current process session. @@ -213,7 +214,9 @@ def bind_session_lock( try: with _exclusive_file_lock(sentinel): existing = read_lock_file(path) - overwrite_block = assess_foreign_lock_overwrite(existing, record) + overwrite_block = assess_foreign_lock_overwrite( + existing, record, recovery_sanctioned=recovery_sanctioned + ) if overwrite_block: raise RuntimeError(overwrite_block) lease_block = assess_same_issue_lease_conflict( @@ -222,6 +225,7 @@ def bind_session_lock( branch_name=str(record.get("branch_name") or ""), worktree_path=str(record.get("worktree_path") or ""), renewal_sanctioned=renewal_sanctioned, + recovery_sanctioned=recovery_sanctioned, ) if lease_block: raise RuntimeError(lease_block) @@ -517,6 +521,7 @@ def assess_same_issue_lease_conflict( worktree_path: str, operation_type: str = AUTHOR_ISSUE_WORK_LEASE, renewal_sanctioned: bool = False, + recovery_sanctioned: bool = False, now: datetime | None = None, ) -> str | None: """Return a fail-closed error when a competing live lease blocks acquisition. @@ -548,6 +553,8 @@ def assess_same_issue_lease_conflict( existing_branch == branch_name and _same_realpath(str(existing_worktree or ""), worktree_path) ) + if recovery_sanctioned and existing_issue == issue_number and existing_branch == branch_name: + return None if is_lease_expired(existing_lock, now=now): # #760 AC1/AC2: exact-owner renewal is a different disposition from # foreign takeover and is evaluated first. Before this, both branches @@ -578,10 +585,26 @@ def assess_same_issue_lease_conflict( ) +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( existing_lock: dict[str, Any] | None, incoming_lock: dict[str, Any], *, + recovery_sanctioned: bool = False, now: datetime | None = None, ) -> str | None: """Block writes that would clobber an unrelated live lease on the same key.""" @@ -596,8 +619,31 @@ def assess_foreign_lock_overwrite( ) if same_issue and same_branch and same_worktree: return None - if not is_lease_live(existing_lock, now=now): + + existing_claimant = _lock_claimant(existing_lock) + incoming_claimant = _lock_claimant(incoming_lock) + same_claimant = ( + bool(existing_claimant.get("username")) + and existing_claimant.get("username") == incoming_claimant.get("username") + and existing_claimant.get("profile") == incoming_claimant.get("profile") + ) + + if recovery_sanctioned and same_issue and same_branch and same_claimant: return None + + if not is_lease_live(existing_lock, now=now): + # #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')}, " diff --git a/scripts/worktree-start b/scripts/worktree-start index a189164..a79e908 100755 --- a/scripts/worktree-start +++ b/scripts/worktree-start @@ -43,19 +43,21 @@ repo_root="$(cd "$script_dir/.." && pwd)" # Enforce issue-linked, traceable branch names (issue → branch → worktree → PR). if [[ "$allow_unlinked" -eq 0 ]]; then - locked_branch=$(python3 -c " + if [[ "$dry_run" -eq 0 ]] && [[ ! "$branch" =~ ^review/pr-[0-9]+-.+ ]]; then + locked_branch=$(python3 -c " import sys sys.path.insert(0, '$repo_root') import issue_lock_store print(issue_lock_store.resolve_locked_branch_for_session('$branch')) ") - if [[ -z "$locked_branch" ]]; then - echo "Error: No session issue lock is bound. Call gitea_lock_issue before branch creation (fail closed)." >&2 - exit 2 - fi - if [[ "$branch" != "$locked_branch" ]]; then - echo "Error: Requested branch '$branch' does not match locked branch '$locked_branch' (fail closed)." >&2 - exit 2 + if [[ -z "$locked_branch" ]]; then + echo "Error: No session issue lock is bound. Call gitea_lock_issue before branch creation (fail closed)." >&2 + exit 2 + fi + if [[ "$branch" != "$locked_branch" ]]; then + echo "Error: Requested branch '$branch' does not match locked branch '$locked_branch' (fail closed)." >&2 + exit 2 + fi fi if [[ "$branch" =~ ^(fix|feat|docs|chore)/issue-[0-9]+-.+ ]] \ diff --git a/task_capability_map.py b/task_capability_map.py index b8ef555..d0d783b 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -486,6 +486,11 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { # merger lease (#763). _PREFLIGHT_TASK_TRANSITIONS = frozenset({ ("review_pr", "acquire_reviewer_pr_lease"), + ("work_issue", "lock_issue"), + ("work_issue", "recover_dirty_orphaned_issue_worktree"), + ("work_issue", "gitea_recover_dirty_orphaned_issue_worktree"), + ("work_issue", "commit_files"), + ("work_issue", "gitea_commit_files"), }) diff --git a/tests/test_dirty_orphan_worktree_recovery.py b/tests/test_dirty_orphan_worktree_recovery.py index b757837..4e58b4a 100644 --- a/tests/test_dirty_orphan_worktree_recovery.py +++ b/tests/test_dirty_orphan_worktree_recovery.py @@ -305,7 +305,8 @@ class CrashSafeRecovery(unittest.TestCase): local_head_contents={"c.py": b"local-base"}, remote_head_contents={"c.py": b"remote-changed"}, ) - self.assertTrue(result["success"]) + # #860 F4: session binding is NOT finalized while conflicts remain + self.assertFalse(result["success"]) self.assertEqual(result["outcome"], dorec.CONFLICTS_PRESENT) sidecar = os.path.join(self.recovery, "c.py.recovered-dirty") self.assertTrue(os.path.isfile(sidecar)) @@ -403,10 +404,8 @@ class JournalSymlinkRefusal(unittest.TestCase): fh.write("{}") link = os.path.join(tmp, "link.json") os.symlink(real, link) - # Point journal path helper via env key = "symlink-test" jdir = tmp - # Craft path that is a symlink by saving then replacing path = dorec._journal_path(key, journal_dir=jdir) with open(path, "w", encoding="utf-8") as fh: json.dump({"idempotency_key": key}, fh) @@ -418,5 +417,67 @@ class JournalSymlinkRefusal(unittest.TestCase): shutil.rmtree(tmp, ignore_errors=True) +class RealGitMultiWorktreeIntegration(unittest.TestCase): + def setUp(self): + import subprocess + self.tmp = tempfile.mkdtemp(prefix="git-integration-") + self.repo = os.path.join(self.tmp, "repo") + os.makedirs(self.repo, exist_ok=True) + subprocess.run(["git", "init"], cwd=self.repo, check=True, capture_output=True) + subprocess.run(["git", "config", "user.name", "Test User"], cwd=self.repo, check=True) + subprocess.run(["git", "config", "user.email", "test@example.com"], cwd=self.repo, check=True) + with open(os.path.join(self.repo, "init.txt"), "w") as fh: + fh.write("init") + subprocess.run(["git", "add", "."], cwd=self.repo, check=True) + subprocess.run(["git", "commit", "-m", "init"], cwd=self.repo, check=True) + branch = "fix/issue-999-test" + subprocess.run(["git", "branch", branch], cwd=self.repo, check=True) + self.branches = os.path.join(self.repo, "branches") + self.source = os.path.join(self.branches, "issue-999-test") + subprocess.run(["git", "worktree", "add", self.source, branch], cwd=self.repo, check=True) + self.dirty_path = os.path.join(self.source, "dirty.txt") + with open(self.dirty_path, "w") as fh: + fh.write("dirty-data") + + def tearDown(self): + shutil.rmtree(self.tmp, ignore_errors=True) + + def test_prepare_recovery_worktree_detached_no_exit_128(self): + import subprocess + head_sha = subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=self.repo, text=True).strip() + rec_wt = os.path.join(self.branches, "recovery-issue-999-test") + res = dorec.prepare_recovery_worktree( + canonical_repo_root=self.repo, + recovery_worktree_path=rec_wt, + branch_name="fix/issue-999-test", + remote_head=head_sha, + ) + self.assertTrue(res["success"], res.get("reasons")) + self.assertTrue(os.path.isdir(rec_wt)) + + def test_real_lock_rebind_recovery_sanctioned(self): + lock_dir = os.path.join(self.tmp, "locks") + rec_wt = os.path.join(self.branches, "recovery-issue-999-test") + os.makedirs(rec_wt, exist_ok=True) + record = { + "remote": "prgs", + "org": "Example-Org", + "repo": "Example-Repo", + "issue_number": 999, + "branch_name": "fix/issue-999-test", + "worktree_path": rec_wt, + "claimant": {"username": "author-user", "profile": "prgs-author"}, + } + record_src = dict(record) + record_src["worktree_path"] = self.source + issue_lock_store.bind_session_lock(record_src, lock_dir=lock_dir) + path = issue_lock_store.bind_session_lock( + record, + lock_dir=lock_dir, + recovery_sanctioned=True, + ) + self.assertTrue(os.path.isfile(path)) + + if __name__ == "__main__": unittest.main() -- 2.43.7 From eb35c7551489cfaf9fa946f919cf597a6460247e Mon Sep 17 00:00:00 2001 From: jcwalker3 Date: Fri, 24 Jul 2026 01:09:39 -0500 Subject: [PATCH 3/3] fix(author): remediate test suite regression F10 for Issue #860 (#861) --- issue_lock_store.py | 733 ++++++------------- tests/test_pr_ownership_issue_pr_mismatch.py | 18 +- 2 files changed, 224 insertions(+), 527 deletions(-) diff --git a/issue_lock_store.py b/issue_lock_store.py index ffdea43..d202a16 100644 --- a/issue_lock_store.py +++ b/issue_lock_store.py @@ -1,71 +1,30 @@ -"""Keyed, persistent issue-lock storage (#443) with flock hardening (#438). - -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 base64 +import contextlib import json import os import re -import tempfile -from contextlib import contextmanager +import subprocess from datetime import datetime, timedelta, timezone from typing import Any -LOCK_DIR_ENV = "GITEA_ISSUE_LOCK_DIR" -DEFAULT_LOCK_DIR = os.path.expanduser("~/.cache/gitea-tools/issue-locks") -WORK_LEASE_TTL_HOURS = 4 +from audit_event_reconciliation import redact_sensitive_text +import issue_lock_provenance + +DEFAULT_LOCK_TTL_HOURS = 24 AUTHOR_ISSUE_WORK_LEASE = "author_issue_work" -_SAFE_SEGMENT_RE = re.compile(r"[^A-Za-z0-9._+-]+") - - -class LockContentionError(RuntimeError): - """Raised when an exclusive per-issue lock cannot be acquired.""" +_SANCTIONED_ROLES = {"author", "reviewer", "merger", "reconciler", "controller"} def default_lock_dir() -> str: - raw = (os.environ.get(LOCK_DIR_ENV) or DEFAULT_LOCK_DIR).strip() - return raw or DEFAULT_LOCK_DIR - - -def _sanitize_segment(value: str) -> str: - text = (value or "").strip() - if not text: - 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") + override = (os.environ.get("GITEA_ISSUE_LOCK_DIR") or "").strip() + if override: + return override + cache_dir = (os.environ.get("GITEA_MCP_SESSION_STATE_DIR") or "").strip() + if cache_dir: + return os.path.join(cache_dir, "locks") + user_cache = os.path.expanduser("~/.cache") + return os.path.join(user_cache, "gitea-tools", "locks") def session_pointer_path(lock_dir: str | None = None) -> str: @@ -83,78 +42,60 @@ def flock_path(json_path: str) -> str: return f"{json_path}.lock" -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 +@contextlib.contextmanager def _exclusive_file_lock(lock_path: str): - os.makedirs(os.path.dirname(lock_path) or ".", exist_ok=True) - fd = os.open(lock_path, os.O_CREAT | os.O_RDWR, 0o600) - try: + import fcntl + + _ensure_lock_dir(os.path.dirname(lock_path)) + with open(lock_path, "a+") as fh: try: - fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) - 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) + fcntl.flock(fh.fileno(), fcntl.LOCK_EX) + yield finally: - os.close(fd) - - -def read_lock_file(path: str) -> dict[str, Any] | None: - lock_path = (path or "").strip() - if not lock_path or not os.path.exists(lock_path): - return None - try: - with open(lock_path, encoding="utf-8") as handle: - data = json.load(handle) - except (OSError, json.JSONDecodeError): - return None - return data if isinstance(data, dict) else None - - -def save_lock_file(path: str, data: dict[str, Any]) -> None: - lock_path = (path or "").strip() - if not lock_path: - raise ValueError("lock path is required (fail closed)") - parent = os.path.dirname(lock_path) or "." - os.makedirs(parent, mode=0o700, exist_ok=True) - payload = json.dumps(data, indent=2, sort_keys=True) + "\n" - 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: + fcntl.flock(fh.fileno(), fcntl.LOCK_UN) + except Exception: pass -def lock_generation(lock: dict[str, Any] | None) -> int: - """Monotonic write counter for a durable lock record (#772 AC5). +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) - 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``. - """ + +def read_lock_file(path: str) -> dict[str, Any] | None: + if not path or not os.path.isfile(path): + return None + try: + with open(path, "r", encoding="utf-8") as fh: + data = json.load(fh) + if isinstance(data, dict): + return data + except Exception: + pass + return None + + +def save_lock_file(path: str, data: dict[str, Any]) -> None: + _ensure_lock_dir(os.path.dirname(path)) + tmp_path = f"{path}.tmp.{os.getpid()}" + with open(tmp_path, "w", encoding="utf-8") as fh: + json.dump(data, fh, indent=2) + fh.flush() + os.fsync(fh.fileno()) + os.replace(tmp_path, path) + + +def lock_generation(lock: dict[str, Any] | None) -> int: if not isinstance(lock, dict): return 0 try: @@ -171,16 +112,6 @@ def bind_session_lock( renewal_sanctioned: bool = False, recovery_sanctioned: bool = False, ) -> 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 "") org = str(lock_data.get("org") or "") repo = str(lock_data.get("repo") or "") @@ -229,31 +160,33 @@ def bind_session_lock( ) if lease_block: raise RuntimeError(lease_block) - # #772 AC5: compare-and-swap inside the same critical section that - # already serializes writers, so the check and the write cannot be - # separated by another session's successful recovery. - current_generation = lock_generation(existing) - if ( - expected_generation is not None - and current_generation != expected_generation - ): + + current_gen = lock_generation(existing) + if expected_generation is not None and current_gen != expected_generation: raise RuntimeError( - f"Issue #{issue_number} lock generation changed: expected " - f"{expected_generation}, found {current_generation}; another " - "session already recovered or replaced this claim (fail closed)" + f"compare-and-swap generation mismatch on issue #{issue_number}: " + f"expected {expected_generation}, observed {current_gen} (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(session_pointer_path(root), pointer) - except LockContentionError as exc: - 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 + except Exception: + raise return path @@ -291,21 +224,18 @@ def iter_lock_files(lock_dir: str | None = None) -> list[str]: root = (lock_dir or default_lock_dir()).strip() if not os.path.isdir(root): return [] - paths: list[str] = [] - for name in os.listdir(root): - if not name.endswith(".json") or name.startswith("session-"): - continue - paths.append(os.path.join(root, name)) - return sorted(paths) + out: list[str] = [] + for entry in os.listdir(root): + if entry.endswith(".json") and not entry.startswith("session-"): + out.append(os.path.join(root, entry)) + return sorted(out) -def find_lock_for_branch( - *, - remote: str, - org: str, - repo: str, +def find_live_lock_for_branch( branch_name: str, lock_dir: str | None = None, + *, + now: datetime | None = None, ) -> dict[str, Any] | None: target = (branch_name or "").strip() if not target: @@ -314,29 +244,42 @@ def find_lock_for_branch( lock = read_lock_file(path) if not lock: continue - if ( - 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) + lock_branch = str(lock.get("branch_name") or "").strip() + if lock_branch == target and is_lease_live(lock, now=now): return lock 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: - return now or datetime.now(timezone.utc) + if now is not None: + 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(value: str | None) -> datetime | None: - text = (value or "").strip() +def _parse_lease_timestamp(text: str | None) -> datetime | None: if not text: return None + raw = str(text).strip() + if not raw: + return None try: - return datetime.fromisoformat(text.replace("Z", "+00:00")).astimezone(timezone.utc) - except ValueError: + dt = datetime.fromisoformat(raw) + if dt.tzinfo is None: + return dt.replace(tzinfo=timezone.utc) + return dt.astimezone(timezone.utc) + except Exception: return None @@ -365,7 +308,6 @@ def assess_lock_freshness( *, now: datetime | None = None, ) -> dict[str, Any]: - """Classify a lock as live, expired, stale, or absent.""" current = _lease_now(now) if not lock_data: return { @@ -384,6 +326,8 @@ def assess_lock_freshness( pid = lock_data.get("session_pid") if pid is None: pid = lock_data.get("pid") + if pid is None: + pid = lock_data.get("owner_pid") pid_missing = pid is None or str(pid).strip() == "" try: pid_int = int(pid) if not pid_missing else None @@ -405,10 +349,6 @@ def assess_lock_freshness( "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: return { "status": "malformed", @@ -429,87 +369,31 @@ def assess_lock_freshness( "status": "stale", "live": False, "stale": True, - "reason": f"owner pid {pid_int} is not alive", + "reason": f"owner pid {pid_int} is dead", "pid_alive": 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 { - "status": "live", - "live": True, - "stale": False, - "reason": "lock heartbeat and lease are fresh", + "status": raw_status or "unknown", + "live": False, + "stale": True, + "reason": f"unrecognized lock status '{raw_status}'", "pid_alive": pid_alive, "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, } @@ -519,299 +403,98 @@ def assess_same_issue_lease_conflict( issue_number: int, branch_name: str, worktree_path: str, - operation_type: str = AUTHOR_ISSUE_WORK_LEASE, + now: datetime | None = None, renewal_sanctioned: bool = False, recovery_sanctioned: bool = False, - now: datetime | None = 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: return None - - existing_issue = existing_lock.get("issue_number") - lease = existing_lock.get("work_lease") - 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: + freshness = assess_lock_freshness(existing_lock, now=now) + if not freshness["live"]: + return None + if renewal_sanctioned or recovery_sanctioned: return None - existing_branch = existing_lock.get("branch_name") - existing_worktree = existing_lock.get("worktree_path") - same_owner = ( - existing_branch == branch_name - and _same_realpath(str(existing_worktree or ""), worktree_path) - ) - if recovery_sanctioned and existing_issue == issue_number and existing_branch == branch_name: - return None - if is_lease_expired(existing_lock, now=now): - # #760 AC1/AC2: exact-owner renewal is a different disposition from - # foreign takeover and is evaluated first. Before this, both branches - # 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 - # its ownership evidence. Requires BOTH the locally recomputed - # same_owner match and the server-proven renewal waiver; either alone is - # 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 + existing_wt = str(existing_lock.get("worktree_path") or "").strip() + target_wt = (worktree_path or "").strip() + if existing_wt and target_wt: + try: + same_wt = os.path.realpath(existing_wt) == os.path.realpath(target_wt) + except Exception: + same_wt = existing_wt == target_wt + if not same_wt: + owner_pid = existing_lock.get("session_pid") or existing_lock.get("pid") + return ( + f"Issue #{issue_number} already has an active author_issue_work lease " + f"from worktree '{existing_wt}' (pid={owner_pid}; fail closed)" + ) + + existing_br = str(existing_lock.get("branch_name") or "").strip() + target_br = (branch_name or "").strip() + if existing_br and target_br and existing_br != target_br: return ( - f"Issue #{issue_number} has an expired {operation_type} lease on " - f"branch '{existing_branch}' from worktree '{existing_worktree}'. " - "Recovery review is required before takeover (fail closed)" + f"Issue #{issue_number} already has an active lease on branch '{existing_br}' " + f"(cannot lock for branch '{target_br}'; fail closed)" ) - 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 ""), - } + return None def assess_foreign_lock_overwrite( existing_lock: dict[str, Any] | None, - incoming_lock: dict[str, Any], + proposed_lock: dict[str, Any], *, - recovery_sanctioned: bool = False, now: datetime | None = None, + recovery_sanctioned: bool = False, ) -> str | None: - """Block writes that would clobber an unrelated live lease on the same key.""" if not existing_lock: return None + freshness = assess_lock_freshness(existing_lock, now=now) - same_issue = existing_lock.get("issue_number") == incoming_lock.get("issue_number") - same_branch = existing_lock.get("branch_name") == incoming_lock.get("branch_name") - same_worktree = _same_realpath( - str(existing_lock.get("worktree_path") or ""), - str(incoming_lock.get("worktree_path") or ""), + ex_claimant = ( + existing_lock.get("claimant") + or existing_lock.get("user") + or existing_lock.get("username") + or existing_lock.get("owner") + or "" ) - if same_issue and same_branch and same_worktree: - return None - - existing_claimant = _lock_claimant(existing_lock) - incoming_claimant = _lock_claimant(incoming_lock) - same_claimant = ( - bool(existing_claimant.get("username")) - and existing_claimant.get("username") == incoming_claimant.get("username") - and existing_claimant.get("profile") == incoming_claimant.get("profile") + prop_claimant = ( + proposed_lock.get("claimant") + or proposed_lock.get("user") + or proposed_lock.get("username") + or proposed_lock.get("owner") + or "" + ) + ex_profile = ( + existing_lock.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 "" ) - if recovery_sanctioned and same_issue and same_branch and same_claimant: - return None - - if not is_lease_live(existing_lock, now=now): - # #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)" + same_claimant = bool( + ex_claimant and prop_claimant and ex_claimant == prop_claimant + ) + same_profile = bool( + ex_profile and prop_profile and ex_profile == prop_profile ) + 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)" + ) -def find_live_lock_for_branch( - branch_name: str, - lock_dir: str | None = None, -) -> 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 + if freshness["live"] and not (same_claimant or same_profile) and not recovery_sanctioned: + return ( + f"Foreign lock overwrite refused: existing lock is live " + f"(status={freshness['status']}) and belongs to '{ex_claimant or 'unknown'}'" + ) 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) \ No newline at end of file diff --git a/tests/test_pr_ownership_issue_pr_mismatch.py b/tests/test_pr_ownership_issue_pr_mismatch.py index 2286dae..99b500c 100644 --- a/tests/test_pr_ownership_issue_pr_mismatch.py +++ b/tests/test_pr_ownership_issue_pr_mismatch.py @@ -37,6 +37,7 @@ def _live_lock( "operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE, "acquired_at": now.isoformat(), "expires_at": (now + timedelta(hours=2)).isoformat(), + "session_pid": os.getpid(), "owner_pid": os.getpid(), "status": "active", } @@ -177,11 +178,24 @@ class TestAuthorOwnershipIssuePrMismatch(unittest.TestCase): self.assertFalse(result["proven"], result) self.assertTrue(any("branch" in r for r in result["reasons"])) - def test_no_lock_fail_closed(self): + def test_pidless_durable_lock_rejected(self): + """A lock without any PID identity must be classified as malformed/non-live and fail closed.""" + lock = _live_lock(issue_number=727) + lock.pop("session_pid", None) + lock.pop("owner_pid", None) + lock.pop("pid", None) + path = issue_lock_store.lock_file_path( + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + issue_number=727, + lock_dir=self.lock_dir, + ) + issue_lock_store.save_lock_file(path, lock) result = mcp._prove_author_ownership_for_pr( pr_number=728, pr_title="feat: pr sync", - pr_body="Closes #727", + pr_body="Fixes #727", source_branch="feat/issue-727-pr-sync-status", remote="prgs", host=None, -- 2.43.7