diff --git a/branch_cleanup_guard.py b/branch_cleanup_guard.py index 64db84a..e91ea02 100644 --- a/branch_cleanup_guard.py +++ b/branch_cleanup_guard.py @@ -525,6 +525,90 @@ def assess_ownership_record_activity(record: dict[str, Any]) -> dict[str, Any]: } +# Reviewer-lease reclaim is only reachable from a non-live (expired/stale) lease. +_RECLAIMABLE_REVIEWER_STATUSES = _EXPIRED_STATUSES | _STALE_STATUSES + + +def is_active_ownership_status(status: str | None) -> bool: + """True when *status* denotes live/active ownership of a branch (#855). + + Used to decide whether a *competing* active claimant still uses a branch + when weighing an expired reviewer lease for reclaim. Expired, stale, + released, and terminal statuses are not active. + """ + return _norm_str(status).lower() in _ACTIVE_OWNERSHIP_STATUSES + + +def assess_expired_reviewer_lease_reclaim( + *, + role: str, + status: str, + pr_merged: bool | None, + owner_pid_alive: bool | None, + competing_active_claimant: bool | None, +) -> dict[str, Any]: + """Decide, explicitly and fail-closed, whether an expired reviewer lease + may stop protecting an already-merged branch (#855 AC4). + + An expired reviewer lease should not protect a merged branch forever once + its work is done and no live claimant remains. Reclaim is permitted only + when **every** condition below is provably satisfied; any unknown + (``None``) or contrary value keeps the lease protective: + + - the lease is a ``reviewer`` lease (author/merger/controller/reconciler + leases are out of scope and always keep protecting); + - its status is expired or stale (never an active/live lease); + - the PR is proven merged (``pr_merged is True``); + - the lease owner process is proven dead (``owner_pid_alive is False``); + - no competing active claimant uses the branch + (``competing_active_claimant is False``). + + Returns a decision dict with ``reclaim_allowed`` and, when refused, the + fail-closed ``reasons``. The reasons never contain secrets — only the + role, the status, and which condition was unproven. + """ + reasons: list[str] = [] + normalized_role = _norm_str(role).lower() + normalized_status = _norm_str(status).lower() + + if normalized_role != "reviewer": + reasons.append( + f"lease role '{normalized_role or 'unknown'}' is not a reviewer " + "lease; expired-reviewer reclaim does not apply" + ) + if normalized_status not in _RECLAIMABLE_REVIEWER_STATUSES: + reasons.append( + f"lease status '{normalized_status or 'unknown'}' is not expired " + "or stale; only a non-live reviewer lease may be reclaimed" + ) + if pr_merged is not True: + reasons.append( + "PR merged state is not proven true; reclaim requires an " + "already-merged PR (fail closed)" + ) + if owner_pid_alive is not False: + reasons.append( + "lease owner process liveness is not proven dead; a live owner " + "still protects the branch (fail closed)" + ) + if competing_active_claimant is not False: + reasons.append( + "a competing active claimant may still use the branch; reclaim " + "requires no other active ownership (fail closed)" + ) + + allowed = not reasons + return { + "reclaim_allowed": allowed, + "role": normalized_role, + "status": normalized_status, + "decision": ( + "reclaim_expired_reviewer_lease" if allowed else "keep_protecting" + ), + "reasons": [] if allowed else reasons, + } + + def assess_active_branch_ownership( *, remote: str, diff --git a/dirty_orphan_worktree_recovery.py b/dirty_orphan_worktree_recovery.py new file mode 100644 index 0000000..17db425 --- /dev/null +++ b/dirty_orphan_worktree_recovery.py @@ -0,0 +1,1080 @@ +"""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: + return { + "success": False, + "created": False, + "resumed": True, + "head": head, + "reasons": [ + f"resume recovery worktree HEAD '{head}' does not match expected remote HEAD '{remote_head}'" + ], + } + return { + "success": True, + "created": False, + "resumed": True, + "head": head, + "reasons": [], + } + + # 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", + "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()}" + ], + } + 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) + + # #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( + 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, + recovery_sanctioned=True, + ) + 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) + + return { + "success": True, + "performed": True, + "outcome": RECOVERY_COMPLETED, + "reasons": [], + "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/docs/mcp-restart-path-inventory.md b/docs/mcp-restart-path-inventory.md new file mode 100644 index 0000000..c880269 --- /dev/null +++ b/docs/mcp-restart-path-inventory.md @@ -0,0 +1,90 @@ +# MCP restart / reload / kill path inventory (#657) + +Complete inventory of every code, script, and host path that can **restart, +reload, reconnect, kill, or force-recreate** an MCP process in this project, +with each path classified and linked to the guard that constrains it. + +This document is the human-readable companion to the machine-readable registry +in [`mcp_restart_paths.py`](../mcp_restart_paths.py). The two are kept in +lock-step by [`tests/test_mcp_restart_paths.py`](../tests/test_mcp_restart_paths.py): +every `path_id` below must appear in this file, and the source guards are run +against the live tree. + +Roadmap linkage: this inventory is the enumeration step of the restart +governance work — parent **#655**, restart-governance ADR **#656**, vision +**#652**, roadmap **#653**. Related detection/guard work: master-advance +staleness **#591**/**#420**, side-effect-free resolver **#685**, transport flap +**#584**, manual-kill contamination **#630**. + +## Classifications + +| Classification | Meaning | +|---|---| +| `sanctioned_narrow_recovery` | One-shot, safe-by-construction recovery that never targets the running daemon. | +| `guarded_fail_closed` | Detects a restart-requiring condition, then fails mutations closed and emits reconnect guidance. Never self-restarts. | +| `forbidden` | A workflow-safety violation; where an LLM tool could invoke it, it is marked contamination. | +| `removed` | A previously-existing unguarded restart primitive that has been deleted; a regression guard keeps it absent. | +| `host_residual` | Behavior owned by the host/IDE, outside this process's control. Documented, not code-guarded here. | + +## The rule + +**No component may perform an unguarded full restart of the MCP daemon.** The +in-process daemon (`gitea_mcp_server.py`, `mcp_server.py`, +`role_session_router.py`) must never replace or terminate its own process: +replacing the process after the host has wired up the stdio pipes desyncs the +JSON-RPC transport (observed with Antigravity/Cascade hosts). Recovery is owned +by the host/operator via a client reconnect — the daemon only ever *detects* +and *fails closed*. + +## Inventory + +| path_id | Classification | Mechanism | Guard | Refs | +|---|---|---|---|---| +| `cli_venv_bootstrap_execv` | sanctioned_narrow_recovery | CLI wrapper scripts re-exec into `venv/bin/python3` via `os.execv`, guarded by `sys.executable != venv_python`. | One-shot pre-import bootstrap; runs before any MCP transport exists and only when not already on the venv interpreter; idempotent guard prevents a re-exec loop. | #657 | +| `daemon_self_replacement` | forbidden | The daemon replacing/terminating its own process (`os.execv`/`os.kill`/`os._exit`) to reload code. | Forbidden by design; enforced against the source tree by `assert_no_daemon_self_replacement()`. | #657, #584 | +| `legacy_auto_restart_helper` | removed | A helper (`_trigger_mcp_auto_restart`) that actively restarted the server from the read-only resolver path. | Removed in #685; kept absent by `assert_auto_restart_helper_absent()`. | #685, #657 | +| `config_touch_reload` | removed | Touching (utime) the MCP client config to make the host reload the server. | Removed from the resolver in #685: stale detection is report-only, never mutating config, spawning threads, or calling `os._exit`. | #685, #657 | +| `master_advance_auto_restart` | guarded_fail_closed | On-disk master advancing past the running code. | `master_parity_gate` captures startup parity and blocks mutations while stale, emitting restart guidance; the process never self-restarts. | #420, #591, #657 | +| `stale_runtime_resolver_reconnect` | guarded_fail_closed | The capability resolver detecting a stale serving process. | Report-only (#685): returns `restart_required`/`stop_required` and an exact reconnect action; no restart, thread, config touch, or `os._exit`. | #685, #657 | +| `manual_daemon_kill` | forbidden | Shell kills of the daemon: `pkill -f mcp_server.py`, `killall`, broad `pkill -f python` sweeps, or `kill ` of a daemon pid. | Forbidden (#630): `runtime_recovery_guard` classifies these as contamination and `gitea_record_daemon_process_kill_attempt` writes a durable marker that fails later mutations closed. Operator maintenance authorization is read only from the environment. | #630, #657 | +| `conflict_marker_infra_stop` | guarded_fail_closed | The daemon entrypoint scans for unresolved merge-conflict markers at startup and stops (`sys.exit(1)`). | Fail-closed startup stop, not a restart: the process exits and waits for the operator to resolve conflicts and relaunch; never loops. | #657 | +| `ide_client_reconnect` | host_residual | A manual `/mcp reconnect` (or equivalent host action) that recreates the MCP client connection. | Outside this process's control; the sanctioned recovery the gates point operators toward. No in-process code initiates it. | #584, #656, #657 | +| `profile_switch_runtime` | sanctioned_narrow_recovery | Switching the active execution profile at runtime (dynamic-profile mode). | In-process and restart-free: `runtime_switching_supported` is true, so a switch rebinds capability without recreating the process. | #656, #657 | + +## Guards enforced in CI + +`tests/test_mcp_restart_paths.py` asserts, against the live source tree: + +1. **Registry well-formedness** — every path has a valid classification, a + non-empty guard description, references, and locations; ids are unique; all + five classifications are represented. +2. **Unknown restart attempts fail closed** — + `assert_restart_attempt_registered()` raises `UnknownRestartPathError` for + any path id not in this inventory, so a novel/unnamed restart primitive + cannot slip through silently. +3. **Daemon never self-replaces** — `assert_no_daemon_self_replacement()` scans + the daemon modules for `os.execv`/`os.kill`/`os._exit`/`os.abort` calls + (comment/docstring mentions are ignored) and finds none. +4. **Legacy helper stays removed** — `assert_auto_restart_helper_absent()` + confirms `_trigger_mcp_auto_restart` has not returned. +5. **pkill stays forbidden** — a daemon `pkill` command still classifies as + contamination via `runtime_recovery_guard`. + +## Residual host behaviors (outside process control) + +* `/mcp reconnect` in the IDE/host — the sanctioned recovery for stale-runtime, + transport-flap (#584), and worktree-binding conditions. The daemon can only + emit guidance toward it. +* Host-level process management (the operator relaunching the daemon after a + fail-closed stop, or after resolving merge conflicts). + +These are documented rather than code-guarded because the process cannot +observe or gate them from inside itself. + +## Rollout + +Per #657, guards are introduced flag-free as **regression assertions** (they +codify invariants that already hold) before any hard runtime block is layered +on. When the restart coordinator (#655/#656) lands, registered paths gain a +coordinator token/capability check; unregistered attempts already fail closed +today via `assert_restart_attempt_registered()`. diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index fb50976..143fa46 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -1469,7 +1469,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, @@ -2050,6 +2050,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 dirty_same_claimant_session_rebind # noqa: E402 # #864 import stacked_pr_support # noqa: E402 import merge_approval_gate # noqa: E402 @@ -4377,6 +4378,280 @@ def gitea_lock_issue( return result +@mcp.tool() +@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 = "" + + 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 + ) + + 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, + 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_locks, + competing_live_sessions=[], + workflow_lease_active=wf_active, + workflow_lease_expired=wf_expired, + 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_rebind_dirty_same_claimant_author_session( issue_number: int, @@ -4633,6 +4908,8 @@ def gitea_rebind_dirty_same_claimant_author_session( } return result + return result + @mcp.tool() def gitea_assess_work_issue_duplicate( @@ -11131,6 +11408,9 @@ def _collect_branch_ownership_records( """ records: list[dict] = [] inventory_error = False + # #855 AC4: expired/stale reviewer-lease records eligible for an explicit + # reclaim decision, evaluated after the full ownership inventory is built. + reviewer_reclaim_candidates: list[tuple[dict, bool | None]] = [] target_branch = (branch or "").strip() if not target_branch: return {"records": records, "inventory_error": False} @@ -11279,15 +11559,28 @@ def _collect_branch_ownership_records( else: status = freshness_status reclaim_allowed = False - records.append( - _base_rec( - category=category, - status=status, - reclaim_allowed=reclaim_allowed, - role=role, - host=lease_host or host_n or host, - ) + rec = _base_rec( + category=category, + status=status, + reclaim_allowed=reclaim_allowed, + role=role, + host=lease_host or host_n or host, ) + records.append(rec) + # #855 AC4: a reviewer lease that is expired/stale (its owner + # gone) becomes a candidate for an explicit, fail-closed + # reclaim decision made once the full inventory is known. + if ( + role == "reviewer" + and status + in branch_cleanup_guard._RECLAIMABLE_REVIEWER_STATUSES + ): + owner_alive = ( + fr.get("owner_pid_alive") if isinstance(fr, dict) else None + ) + reviewer_reclaim_candidates.append( + (rec, owner_alive if isinstance(owner_alive, bool) else None) + ) except Exception: # O1: fail closed on control-plane inventory errors. inventory_error = True @@ -11358,6 +11651,44 @@ def _collect_branch_ownership_records( ) ) + # #855 AC4: decide, explicitly and fail-closed, whether any expired/stale + # reviewer lease may stop protecting an already-merged branch. This runs + # only after the full ownership inventory is built, so a competing active + # claimant (an active lease, author session, worktree binding, or active + # reviewer comment lease) is visible. An inventory failure keeps every + # reclaim candidate protective (reclaim_allowed stays False). + if reviewer_reclaim_candidates and not inventory_error: + pr_merged_state: bool | None = None + if pr_number is not None and auth and base_api: + try: + pr_live = api_request( + "GET", f"{base_api}/pulls/{int(pr_number)}", auth + ) + if isinstance(pr_live, dict) and pr_live: + pr_merged_state = bool( + pr_live.get("merged") or pr_live.get("merged_at") + ) + except Exception: + # Unknown merged state fails closed (candidate stays protective). + pr_merged_state = None + for cand_rec, owner_alive in reviewer_reclaim_candidates: + competing = any( + other is not cand_rec + and branch_cleanup_guard.is_active_ownership_status( + other.get("status") + ) + for other in records + ) + decision = branch_cleanup_guard.assess_expired_reviewer_lease_reclaim( + role=str(cand_rec.get("role")), + status=str(cand_rec.get("status")), + pr_merged=pr_merged_state, + owner_pid_alive=owner_alive, + competing_active_claimant=competing, + ) + cand_rec["reclaim_allowed"] = decision["reclaim_allowed"] + cand_rec["reclaim_decision"] = decision["decision"] + return {"records": records, "inventory_error": inventory_error} @@ -11420,6 +11751,7 @@ def gitea_reconcile_merged_cleanups( dry_run: bool = True, execute_confirmed: bool = False, limit: int = 50, + pr_number: int | None = None, remote: str = "dadeschools", host: str | None = None, org: str | None = None, @@ -11430,7 +11762,11 @@ def gitea_reconcile_merged_cleanups( Args: dry_run: Defaults to True. When True, only builds the reconciliation report. execute_confirmed: Must be True when dry_run=False. - limit: Max number of closed PRs to inspect. + limit: Max number of closed PRs to inspect (batch mode only; ignored when + ``pr_number`` is set). + pr_number: Optional exact merged PR selector (#855). When set, only that + PR is assessed/acted on (fail closed if missing, unmerged, or + ambiguous). When omitted, existing batch behaviour is preserved. remote: Known Gitea instance ('dadeschools' or 'prgs'). host: Override the Gitea host. org: Override the owner/organization. @@ -11465,11 +11801,120 @@ def gitea_reconcile_merged_cleanups( "audit_phase": audit_reconciliation_mode.current_phase(), } + # #855: optional exact PR pin. Fail closed before any inventory mutation. + exact_pr: int | None = None + if pr_number is not None: + try: + exact_pr = int(pr_number) + except (TypeError, ValueError): + return { + "success": False, + "performed": False, + "executed": False, + "dry_run": bool(dry_run), + "selection_mode": "exact_pr", + "selected_pr_number": pr_number, + "reasons": [ + f"pr_number={pr_number!r} is not a valid integer " + "(fail closed; no mutation)" + ], + "blocker_kind": "invalid_pr_number", + } + if exact_pr <= 0: + return { + "success": False, + "performed": False, + "executed": False, + "dry_run": bool(dry_run), + "selection_mode": "exact_pr", + "selected_pr_number": exact_pr, + "reasons": [ + f"pr_number={exact_pr} must be a positive integer " + "(fail closed; no mutation)" + ], + "blocker_kind": "invalid_pr_number", + } + h, o, r = _resolve(remote, host, org, repo) auth = _auth(h) base = repo_api_url(h, o, r) - closed_prs = api_get_all(f"{base}/pulls?state=closed", auth, limit=limit) - open_prs = api_get_all(f"{base}/pulls?state=open", auth) + + selection_mode = "batch" + closed_prs: list[dict] = [] + open_prs: list[dict] = [] + if exact_pr is not None: + selection_mode = "exact_pr" + try: + pr_live = api_request("GET", f"{base}/pulls/{exact_pr}", auth) + except Exception as exc: + return { + "success": False, + "performed": False, + "executed": False, + "dry_run": bool(dry_run), + "selection_mode": selection_mode, + "selected_pr_number": exact_pr, + "reasons": [ + f"PR #{exact_pr} could not be uniquely resolved " + f"(fail closed; no mutation): {_redact(str(exc))}" + ], + "blocker_kind": "pr_unresolvable", + } + if not isinstance(pr_live, dict) or not pr_live: + return { + "success": False, + "performed": False, + "executed": False, + "dry_run": bool(dry_run), + "selection_mode": selection_mode, + "selected_pr_number": exact_pr, + "reasons": [ + f"PR #{exact_pr} could not be uniquely resolved " + "(empty response; fail closed; no mutation)" + ], + "blocker_kind": "pr_unresolvable", + } + live_number = pr_live.get("number") + try: + live_number_int = int(live_number) if live_number is not None else None + except (TypeError, ValueError): + live_number_int = None + if live_number_int != exact_pr: + return { + "success": False, + "performed": False, + "executed": False, + "dry_run": bool(dry_run), + "selection_mode": selection_mode, + "selected_pr_number": exact_pr, + "reasons": [ + f"PR #{exact_pr} resolution is ambiguous or mismatched " + f"(live number={live_number!r}; fail closed; no mutation)" + ], + "blocker_kind": "pr_ambiguous", + } + if not (pr_live.get("merged") or pr_live.get("merged_at")): + return { + "success": False, + "performed": False, + "executed": False, + "dry_run": bool(dry_run), + "selection_mode": selection_mode, + "selected_pr_number": exact_pr, + "reasons": [ + f"PR #{exact_pr} is not merged " + "(exact-target cleanup requires a merged PR; " + "fail closed; no mutation)" + ], + "blocker_kind": "pr_not_merged", + } + closed_prs = [pr_live] + # Exact mode still needs open heads for remote-delete safety gates. + open_prs = api_get_all(f"{base}/pulls?state=open", auth) + else: + # Preserve historical call order (closed then open) for batch callers/tests. + closed_prs = api_get_all(f"{base}/pulls?state=closed", auth, limit=limit) + open_prs = api_get_all(f"{base}/pulls?state=open", auth) merged_closed: list[dict] = [] remote_branch_exists: dict[str, bool] = {} @@ -11496,6 +11941,13 @@ def gitea_reconcile_merged_cleanups( scratch_candidates = merged_cleanup_reconcile.discover_reviewer_scratch_worktrees( _canonical_local_git_root() ) + # #855: exact-target never inventories or mutates foreign PR scratch trees. + if exact_pr is not None: + scratch_candidates = [ + s + for s in scratch_candidates + if int(s.get("pr_number") or 0) == int(exact_pr) + ] active_reviewer_leases: dict[int, bool] = {} pr_states: dict[int, dict] = {} for scratch in scratch_candidates: @@ -11530,6 +11982,33 @@ def gitea_reconcile_merged_cleanups( active_reviewer_leases=active_reviewer_leases, pr_states=pr_states, ) + report["selection_mode"] = selection_mode + if exact_pr is not None: + report["selected_pr_number"] = exact_pr + # Fail closed if exact pin somehow produced other or zero entries. + entries = list(report.get("entries") or []) + entry_numbers = [] + for entry in entries: + try: + entry_numbers.append(int(entry.get("pr_number"))) + except (TypeError, ValueError): + entry_numbers.append(entry.get("pr_number")) + if entry_numbers != [exact_pr]: + return { + "success": False, + "performed": False, + "executed": False, + "dry_run": bool(dry_run), + "selection_mode": selection_mode, + "selected_pr_number": exact_pr, + "reasons": [ + f"exact PR #{exact_pr} selection produced unexpected " + f"candidate set {entry_numbers!r} " + "(fail closed; no mutation)" + ], + "blocker_kind": "exact_selection_mismatch", + "entries": entries, + } if dry_run: report["dry_run"] = True diff --git a/issue_lock_provenance.py b/issue_lock_provenance.py index 544e017..b79d8ad 100644 --- a/issue_lock_provenance.py +++ b/issue_lock_provenance.py @@ -16,6 +16,7 @@ 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" # #864: dirty-preserving same-claimant author-session rebind (dead owner PID). SOURCE_DIRTY_SAME_CLAIMANT_REBIND = ( "gitea_rebind_dirty_same_claimant_author_session" @@ -25,6 +26,7 @@ SANCTIONED_LOCK_SOURCES = frozenset({ SOURCE_LOCK_ISSUE, SOURCE_LOCK_ADOPTION, SOURCE_OPERATOR_OVERRIDE, + SOURCE_RECOVER_DIRTY_ORPHANED, SOURCE_DIRTY_SAME_CLAIMANT_REBIND, }) diff --git a/issue_lock_store.py b/issue_lock_store.py index a963174..d011c1c 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) @@ -380,7 +384,18 @@ 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 + 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 + 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 +404,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 +442,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, } @@ -486,6 +523,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. @@ -517,6 +555,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 @@ -547,10 +587,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.""" @@ -565,8 +621,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/mcp_restart_paths.py b/mcp_restart_paths.py new file mode 100644 index 0000000..087fe72 --- /dev/null +++ b/mcp_restart_paths.py @@ -0,0 +1,475 @@ +"""Inventory and fail-closed guards for MCP restart/reload/kill paths (#657). + +Single source of truth enumerating every code/script/doc path that can +restart, reload, reconnect, kill, or force-recreate an MCP process. Each path +is classified and linked to the guard that constrains it. The companion +human-readable inventory lives in ``docs/mcp-restart-path-inventory.md`` and is +kept in lock-step with this module by ``tests/test_mcp_restart_paths.py``. + +Design intent (aligns with #655 restart-coordinator roadmap): + +* **No unguarded full restart.** The in-process MCP daemon + (``gitea_mcp_server.py`` / ``mcp_server.py`` / ``role_session_router.py``) + must never replace or kill its own process — replacing the process after the + host wired up the stdio pipes desyncs the JSON-RPC transport (observed with + Antigravity/Cascade hosts). ``assert_no_daemon_self_replacement`` enforces + this against the live source tree. +* **No legacy auto-restart helper.** ``_trigger_mcp_auto_restart`` was removed + when the stale-runtime resolver became side-effect free (#685); + ``assert_auto_restart_helper_absent`` keeps it removed. +* **Unknown restart attempts fail closed.** LLM tools must route any restart + intent through a *registered* path. ``assert_restart_attempt_registered`` + raises ``UnknownRestartPathError`` for anything not in this inventory. +* **pkill stays forbidden (#630).** Manual daemon kills are classified as + contamination by :mod:`runtime_recovery_guard`; this module records that path + and the test asserts the classification still holds. + +This module performs no restarts, spawns no threads, and touches no config or +process state. It is pure inventory + read-only source assertions. +""" + +from __future__ import annotations + +import os +from dataclasses import dataclass +from pathlib import Path +from typing import Iterable + +# --- Classifications ------------------------------------------------------- + +#: A narrow, one-shot recovery that is safe by construction (e.g. a CLI wrapper +#: re-execing into the venv interpreter before importing anything, or an +#: in-process profile switch). Never targets the running MCP daemon process. +CLASS_SANCTIONED_NARROW = "sanctioned_narrow_recovery" + +#: The path detects a condition that would require a restart, then *fails +#: closed* on mutations and emits restart/reconnect guidance. It never restarts +#: the process itself (recovery is owned by the host/operator). +CLASS_GUARDED_FAIL_CLOSED = "guarded_fail_closed" + +#: The path is forbidden. Attempting it is a workflow-safety violation and, +#: where an LLM tool could invoke it, is marked as contamination. +CLASS_FORBIDDEN = "forbidden" + +#: A previously-existing unguarded restart primitive that has been deleted. A +#: regression guard keeps it absent. +CLASS_REMOVED = "removed" + +#: Behavior that lives in the host/IDE and is outside this process's control +#: (e.g. a manual ``/mcp reconnect``). Documented, not code-guarded here. +CLASS_HOST_RESIDUAL = "host_residual" + +VALID_CLASSIFICATIONS = frozenset( + { + CLASS_SANCTIONED_NARROW, + CLASS_GUARDED_FAIL_CLOSED, + CLASS_FORBIDDEN, + CLASS_REMOVED, + CLASS_HOST_RESIDUAL, + } +) + +#: The in-process MCP daemon modules. These must never self-replace/self-kill. +DAEMON_MODULES = ( + "gitea_mcp_server.py", + "mcp_server.py", + "role_session_router.py", +) + +#: The legacy auto-restart helper removed in #685. Must stay removed. +LEGACY_AUTO_RESTART_HELPER = "_trigger_mcp_auto_restart" + +#: Call patterns that would let the daemon replace or terminate its own +#: process. Matched as calls (trailing ``(``) so prose/docstring mentions such +#: as "we do NOT os.execv() here" or "never calls ``os._exit``" do not trip the +#: scanner (comment lines are stripped first regardless). +DAEMON_SELF_REPLACEMENT_PRIMITIVES = ( + "os.execv(", + "os.execve(", + "os.execvp(", + "os.execvpe(", + "os.kill(", + "os.killpg(", + "os._exit(", + "os.abort(", +) + + +@dataclass(frozen=True) +class RestartPath: + """One classified restart/reload/kill path in the inventory.""" + + path_id: str + title: str + mechanism: str + classification: str + guard: str + locations: tuple[str, ...] + references: tuple[str, ...] + residual_host: bool = False + notes: str = "" + + +class UnknownRestartPathError(RuntimeError): + """Raised when a restart attempt is not a registered, classified path.""" + + +# --- The inventory --------------------------------------------------------- + +_RESTART_PATHS: tuple[RestartPath, ...] = ( + RestartPath( + path_id="cli_venv_bootstrap_execv", + title="CLI wrapper venv re-exec", + mechanism=( + "Standalone CLI scripts re-exec into venv/bin/python3 via os.execv " + "at import top, guarded by `sys.executable != venv_python`." + ), + classification=CLASS_SANCTIONED_NARROW, + guard=( + "One-shot, pre-import bootstrap; runs before any MCP transport " + "exists and only when not already on the venv interpreter, so it " + "cannot desync a live daemon. Idempotent guard condition prevents " + "a re-exec loop." + ), + locations=( + "create_pr.py", + "create_issue.py", + "close_issue.py", + "merge_pr.py", + "review_pr.py", + "edit_pr.py", + "delete_branch.py", + "mark_issue.py", + "manage_labels.py", + "list_issues.py", + "list_prs.py", + ), + references=("#657",), + ), + RestartPath( + path_id="daemon_self_replacement", + title="MCP daemon self-replacement", + mechanism=( + "The in-process MCP daemon replacing/terminating its own process " + "(os.execv/os.kill/os._exit) to reload code." + ), + classification=CLASS_FORBIDDEN, + guard=( + "Forbidden by design: replacing the process after the host wired " + "up stdio desyncs JSON-RPC (Antigravity/Cascade). Enforced against " + "the source tree by assert_no_daemon_self_replacement()." + ), + locations=("gitea_mcp_server.py:~155 (decision comment)",) + DAEMON_MODULES, + references=("#657", "#584"), + ), + RestartPath( + path_id="legacy_auto_restart_helper", + title="Legacy _trigger_mcp_auto_restart helper", + mechanism=( + "A helper that actively restarted the MCP server from the " + "read-only resolver path." + ), + classification=CLASS_REMOVED, + guard=( + "Removed in #685 when the resolver became side-effect free. Kept " + "absent by assert_auto_restart_helper_absent()." + ), + locations=("gitea_mcp_server.py", "mcp_server.py"), + references=("#685", "#657"), + ), + RestartPath( + path_id="config_touch_reload", + title="MCP client config-touch reload", + mechanism=( + "Touching (utime) the MCP client config file to make the host " + "reload/recreate the server process." + ), + classification=CLASS_REMOVED, + guard=( + "Removed from the resolver in #685: stale-runtime detection is " + "report-only and never mutates client config, spawns threads, or " + "calls os._exit." + ), + locations=("gitea_mcp_server.py (resolve_task_capability)",), + references=("#685", "#657"), + ), + RestartPath( + path_id="master_advance_auto_restart", + title="Master-advance staleness gate", + mechanism=( + "On-disk master advancing past the running code. The master-parity " + "gate detects it and fails mutations closed with restart guidance." + ), + classification=CLASS_GUARDED_FAIL_CLOSED, + guard=( + "Detect + fail closed only; the process never self-restarts. " + "master_parity_gate captures startup parity and blocks mutations " + "while stale, emitting restart/reconnect guidance." + ), + locations=( + "master_parity_gate.py", + "gitea_mcp_server.py (gitea_assess_master_parity)", + ), + references=("#420", "#591", "#657"), + ), + RestartPath( + path_id="stale_runtime_resolver_reconnect", + title="Stale-runtime resolver reconnect guidance", + mechanism=( + "The capability resolver detecting a stale serving process and " + "reporting restart_required/stop_required for a client reconnect." + ), + classification=CLASS_GUARDED_FAIL_CLOSED, + guard=( + "Report-only (#685): returns restart_required/stop_required and an " + "exact_safe_next_action pointing at IDE/client reconnect; performs " + "no restart, thread spawn, config touch, or os._exit." + ), + locations=("gitea_mcp_server.py (gitea_resolve_task_capability)",), + references=("#685", "#657"), + ), + RestartPath( + path_id="manual_daemon_kill", + title="Manual daemon kill (pkill/killall/kill)", + mechanism=( + "Shell kills of the MCP daemon: `pkill -f mcp_server.py`, " + "`killall`, broad `pkill -f python` sweeps, or `kill ` of a " + "daemon pid." + ), + classification=CLASS_FORBIDDEN, + guard=( + "Forbidden (#630): runtime_recovery_guard classifies these as " + "contamination and gitea_record_daemon_process_kill_attempt writes " + "a durable marker that fails subsequent mutations closed. Operator " + "maintenance authorization is read only from the environment, not " + "from a tool argument." + ), + locations=( + "runtime_recovery_guard.py", + "gitea_mcp_server.py (gitea_record_daemon_process_kill_attempt)", + ), + references=("#630", "#657"), + ), + RestartPath( + path_id="conflict_marker_infra_stop", + title="Startup conflict-marker infra stop", + mechanism=( + "The daemon entrypoint scans for unresolved merge-conflict markers " + "at startup and stops (sys.exit(1)) if found." + ), + classification=CLASS_GUARDED_FAIL_CLOSED, + guard=( + "Fail-closed startup stop, not a restart: the process exits and " + "waits for the operator to resolve conflicts and relaunch. Never " + "self-restarts or loops." + ), + locations=("mcp_server.py (check_conflict_markers)",), + references=("#657",), + ), + RestartPath( + path_id="ide_client_reconnect", + title="Host/IDE MCP reconnect", + mechanism=( + "A manual `/mcp reconnect` (or equivalent host action) that the " + "IDE performs to recreate the MCP client connection." + ), + classification=CLASS_HOST_RESIDUAL, + guard=( + "Outside this process's control. It is the sanctioned recovery the " + "gates point operators toward; documented as residual host " + "behavior. No in-process code initiates it." + ), + locations=("host/IDE",), + references=("#584", "#656", "#657"), + residual_host=True, + ), + RestartPath( + path_id="profile_switch_runtime", + title="Runtime profile switch", + mechanism=( + "Switching the active execution profile at runtime " + "(dynamic-profile mode)." + ), + classification=CLASS_SANCTIONED_NARROW, + guard=( + "In-process and restart-free: runtime_switching_supported is true, " + "so a profile switch rebinds capability without recreating the " + "process. No restart primitive is invoked." + ), + locations=("gitea_mcp_server.py (gitea_activate_profile)",), + references=("#656", "#657"), + ), +) + +_BY_ID: dict[str, RestartPath] = {p.path_id: p for p in _RESTART_PATHS} + + +# --- Read-only accessors --------------------------------------------------- + + +def iter_restart_paths() -> tuple[RestartPath, ...]: + """Return the full inventory as an immutable tuple.""" + + return _RESTART_PATHS + + +def restart_path_ids() -> frozenset[str]: + """Return the set of registered path ids.""" + + return frozenset(_BY_ID) + + +def get_restart_path(path_id: str) -> RestartPath: + """Return the registered path, or raise :class:`UnknownRestartPathError`.""" + + try: + return _BY_ID[path_id] + except KeyError as exc: + raise UnknownRestartPathError( + f"unknown restart path id {path_id!r}; not in the #657 inventory" + ) from exc + + +def paths_by_classification(classification: str) -> tuple[RestartPath, ...]: + """Return all registered paths with the given classification.""" + + if classification not in VALID_CLASSIFICATIONS: + raise ValueError(f"unknown classification {classification!r}") + return tuple(p for p in _RESTART_PATHS if p.classification == classification) + + +def assert_restart_attempt_registered(path_id: str) -> RestartPath: + """Fail closed unless ``path_id`` is a registered, classified restart path. + + LLM tools that intend to trigger any restart/reload/reconnect must name a + registered path so an unknown/novel restart primitive cannot slip through + silently. Forbidden and removed paths are registered too — this only + asserts the attempt is *known*, not that it is *permitted*; callers must + still honor the classification. + """ + + return get_restart_path(path_id) + + +def assert_registry_wellformed() -> None: + """Validate the inventory's own invariants (fail closed on drift).""" + + seen: set[str] = set() + for path in _RESTART_PATHS: + if path.path_id in seen: + raise ValueError(f"duplicate restart path id {path.path_id!r}") + seen.add(path.path_id) + if path.classification not in VALID_CLASSIFICATIONS: + raise ValueError( + f"{path.path_id!r} has invalid classification " + f"{path.classification!r}" + ) + if not path.guard.strip(): + raise ValueError(f"{path.path_id!r} is missing a guard description") + if not path.references: + raise ValueError(f"{path.path_id!r} is missing references") + if not path.locations: + raise ValueError(f"{path.path_id!r} is missing locations") + if path.classification == CLASS_HOST_RESIDUAL and not path.residual_host: + raise ValueError( + f"{path.path_id!r} is host_residual but residual_host is False" + ) + + +# --- Source-tree guards ---------------------------------------------------- + + +def _repo_root(root: str | os.PathLike[str] | None = None) -> Path: + if root is not None: + return Path(root) + return Path(__file__).resolve().parent + + +def _iter_code_lines(text: str) -> Iterable[tuple[int, str]]: + """Yield (1-based lineno, line) for lines that are not full-line comments.""" + + for lineno, line in enumerate(text.splitlines(), start=1): + if line.lstrip().startswith("#"): + continue + yield lineno, line + + +def scan_daemon_self_replacement( + root: str | os.PathLike[str] | None = None, +) -> list[dict[str, object]]: + """Return violations where a daemon module could self-replace/self-kill. + + Scans :data:`DAEMON_MODULES` for calls in + :data:`DAEMON_SELF_REPLACEMENT_PRIMITIVES`. Full-line comments are ignored, + and only call forms (with a trailing ``(``) match, so decision comments and + docstrings that merely mention the primitives do not produce false hits. + """ + + repo = _repo_root(root) + violations: list[dict[str, object]] = [] + for module in DAEMON_MODULES: + path = repo / module + if not path.exists(): + continue + text = path.read_text(encoding="utf-8", errors="replace") + for lineno, line in _iter_code_lines(text): + for primitive in DAEMON_SELF_REPLACEMENT_PRIMITIVES: + if primitive in line: + violations.append( + { + "module": module, + "line": lineno, + "primitive": primitive, + "text": line.strip(), + } + ) + return violations + + +def assert_no_daemon_self_replacement( + root: str | os.PathLike[str] | None = None, +) -> None: + """Fail closed if any daemon module can restart/kill its own process.""" + + violations = scan_daemon_self_replacement(root) + if violations: + rendered = "; ".join( + f"{v['module']}:{v['line']} {v['primitive']}" for v in violations + ) + raise AssertionError( + "MCP daemon must never self-replace/self-kill (#657); found: " + f"{rendered}" + ) + + +def scan_auto_restart_helper( + root: str | os.PathLike[str] | None = None, +) -> list[dict[str, object]]: + """Return occurrences of a *definition* of the legacy auto-restart helper.""" + + repo = _repo_root(root) + needle = f"def {LEGACY_AUTO_RESTART_HELPER}" + hits: list[dict[str, object]] = [] + for module in DAEMON_MODULES: + path = repo / module + if not path.exists(): + continue + text = path.read_text(encoding="utf-8", errors="replace") + for lineno, line in _iter_code_lines(text): + if needle in line: + hits.append({"module": module, "line": lineno}) + return hits + + +def assert_auto_restart_helper_absent( + root: str | os.PathLike[str] | None = None, +) -> None: + """Fail closed if the removed ``_trigger_mcp_auto_restart`` reappears.""" + + hits = scan_auto_restart_helper(root) + if hits: + rendered = "; ".join(f"{h['module']}:{h['line']}" for h in hits) + raise AssertionError( + f"{LEGACY_AUTO_RESTART_HELPER} was removed in #685 and must not " + f"return (#657); found definition at: {rendered}" + ) 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 7d7b76b..213e701 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", + }, # #864: dirty-preserving same-claimant author-session rebind (dead owner PID). # Author MCP tool path. Reconciler execute is gated inside the tool via # authorize_reconciler_execute + role_kind checks (not this map entry). @@ -488,6 +497,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_branch_cleanup_guard.py b/tests/test_branch_cleanup_guard.py index 0bd1fda..fddb534 100644 --- a/tests/test_branch_cleanup_guard.py +++ b/tests/test_branch_cleanup_guard.py @@ -1639,6 +1639,358 @@ class TestSecondRemediationIntegration(unittest.TestCase): self.assertTrue(ownership_calls) +class TestIssue855ExactPrSelector(unittest.TestCase): + """#855: exact pr_number pin for reconcile_merged_cleanups (#851 lifecycle).""" + + def setUp(self): + self._remotes = patch.dict( + mcp_server.REMOTES, + { + "prgs": { + "host": "gitea.example.com", + "org": "Scaled-Tech-Consulting", + "repo": "Gitea-Tools", + } + }, + ) + self._remotes.start() + patch("gitea_audit.audit_enabled", return_value=False).start() + self.mock_api = patch("mcp_server.api_request").start() + self.mock_all = patch("mcp_server.api_get_all", return_value=[]).start() + patch("mcp_server.get_auth_header", return_value=FAKE_AUTH).start() + patch( + "mcp_server.merged_cleanup_reconcile.is_head_ancestor_of_ref", + return_value=True, + ).start() + patch( + "mcp_server.get_profile", + return_value=dict(RECONCILER_WITH_DELETE), + ).start() + patch( + "mcp_server._profile_operation_gate", + return_value=[], + ).start() + patch( + "mcp_server._collect_branch_ownership_records", + return_value={"records": [], "inventory_error": False}, + ).start() + patch( + "mcp_server.merged_cleanup_reconcile.discover_reviewer_scratch_worktrees", + return_value=[], + ).start() + patch("mcp_server.verify_preflight_purity", return_value=None).start() + patch( + "mcp_server.audit_reconciliation_mode.check_cleanup_execution_allowed", + return_value=(True, []), + ).start() + + def tearDown(self): + patch.stopall() + + def _merged_pr(self, number, branch, sha="c" * 40): + return { + "number": number, + "title": f"PR {number}", + "body": f"Closes #{number - 4}", + "merged": True, + "merged_at": "2026-07-23T12:00:00Z", + "merge_commit_sha": "f" * 40, + "state": "closed", + "head": {"ref": branch, "sha": sha}, + "base": {"ref": "master"}, + } + + def test_exact_pr_848_ignores_newer_852_in_batch_queue(self): + """pr_number=848 selects only #848 even when #852 is newer/first.""" + from mcp_server import gitea_reconcile_merged_cleanups + + pr_848 = self._merged_pr( + 848, "fix/issue-844-exclude-epic-containers", sha="c3f282ba" + "0" * 32 + ) + # Closed list would rank #852 first in batch mode; exact pin must ignore it. + closed_batch = [ + self._merged_pr(852, "fix/issue-851-cleanup-worktree-before-remote-delete"), + pr_848, + self._merged_pr(849, "fix/issue-849-other"), + self._merged_pr(846, "fix/issue-846-other"), + self._merged_pr(845, "fix/issue-845-other"), + ] + batch_fetch_calls = [] + + def fake_api(method, url, *args, **kwargs): + if method == "GET" and url.rstrip("/").endswith("/pulls/848"): + return dict(pr_848) + if method == "GET" and "/pulls/" in url: + raise AssertionError(f"unexpected PR fetch: {url}") + if method == "GET" and "/branches/" in url: + return {"name": "present"} + return {} + + def fake_all(url, auth, limit=None): + batch_fetch_calls.append((url, limit)) + if "state=open" in url: + return [] + if "state=closed" in url: + # Exact mode must not use the closed batch list. + raise AssertionError( + "exact pr_number mode must not page closed PRs: " + url + ) + return [] + + self.mock_api.side_effect = fake_api + self.mock_all.side_effect = fake_all + patch( + "mcp_server._remote_branch_exists", + return_value=True, + ).start() + patch( + "mcp_server.merged_cleanup_reconcile.build_reconciliation_report", + side_effect=lambda **kwargs: { + "entries": [ + { + "pr_number": int(pr["number"]), + "head_branch": (pr.get("head") or {}).get("ref"), + "issue_number": 844, + "remote_branch": { + "safe_to_delete_remote": True, + "head_branch": (pr.get("head") or {}).get("ref"), + }, + "local_worktree": { + "safe_to_remove_worktree": True, + "worktree_path": ( + "/tmp/branches/fix-issue-844-exclude-epic-containers" + ), + }, + "planned_execution_order": ( + mcp_server.merged_cleanup_reconcile.plan_cleanup_execution_order( + remote_assessment={"safe_to_delete_remote": True}, + local_assessment={"safe_to_remove_worktree": True}, + ) + ), + } + for pr in kwargs.get("closed_prs") or [] + if pr.get("merged_at") or pr.get("merged") + ], + "reviewer_scratch_entries": [], + "merged_pr_count": len(kwargs.get("closed_prs") or []), + }, + ).start() + + res = gitea_reconcile_merged_cleanups( + dry_run=True, + pr_number=848, + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + ) + self.assertTrue(res.get("success")) + self.assertFalse(res.get("performed")) + self.assertEqual(res.get("selection_mode"), "exact_pr") + self.assertEqual(res.get("selected_pr_number"), 848) + entries = res.get("entries") or [] + self.assertEqual(len(entries), 1, entries) + self.assertEqual(entries[0].get("pr_number"), 848) + self.assertEqual( + entries[0].get("head_branch"), + "fix/issue-844-exclude-epic-containers", + ) + # No other PR appears in plan. + self.assertEqual(list((res.get("planned_execution_orders") or {}).keys()), ["848"]) + plan = (res.get("planned_execution_orders") or {}).get("848") or [] + actions = [s.get("action") for s in plan] + self.assertEqual( + actions, + [ + "remove_local_worktree", + "reassess_branch_ownership", + "delete_remote_branch", + ], + ) + # Prove we never scanned the multi-PR closed batch. + self.assertFalse(any("state=closed" in (u or "") for u, _ in batch_fetch_calls)) + # closed_batch fixture must remain unused (sanity). + self.assertEqual(closed_batch[0]["number"], 852) + + def test_exact_pr_execute_only_mutates_selected_pr(self): + """Execute with pr_number must never touch #845/#846/#849/#852.""" + from mcp_server import gitea_reconcile_merged_cleanups + + pr_848 = self._merged_pr(848, "fix/issue-844-exclude-epic-containers") + worktree_path = "/tmp/branches/fix-issue-844-exclude-epic-containers" + remove_calls = [] + delete_api_calls = [] + ownership_branches = [] + + def fake_api(method, url, *args, **kwargs): + if method == "GET" and url.rstrip("/").endswith("/pulls/848"): + return dict(pr_848) + if method == "DELETE": + delete_api_calls.append(url) + # Forbid foreign PR branch deletion by URL content. + for forbidden in ("845", "846", "849", "852"): + self.assertNotIn(forbidden, url) + return {} + + def fake_remove(project_root, branch, worktree_path=None): + remove_calls.append({"branch": branch, "worktree_path": worktree_path}) + return { + "success": True, + "performed": True, + "message": f"removed {worktree_path}", + "worktree_path": worktree_path, + } + + def fake_collect(**kwargs): + ownership_branches.append(kwargs.get("branch")) + return {"records": [], "inventory_error": False} + + def fake_probe(h, o, r, auth, br): + return guard.classify_branch_readback_http_status( + 404, not_found_scope=guard.NOT_FOUND_SCOPE_BRANCH + ) + + self.mock_api.side_effect = fake_api + self.mock_all.side_effect = lambda url, auth, limit=None: [] + patch("mcp_server._remote_branch_exists", return_value=True).start() + patch( + "mcp_server.merged_cleanup_reconcile.build_reconciliation_report", + return_value={ + "entries": [ + { + "pr_number": 848, + "head_branch": "fix/issue-844-exclude-epic-containers", + "remote_branch": {"safe_to_delete_remote": True}, + "local_worktree": { + "safe_to_remove_worktree": True, + "worktree_path": worktree_path, + }, + "planned_execution_order": [ + {"action": "remove_local_worktree", "phase": 1}, + {"action": "reassess_branch_ownership", "phase": 2}, + {"action": "delete_remote_branch", "phase": 3}, + ], + } + ], + "reviewer_scratch_entries": [ + # Foreign scratch must be filtered before report execute loop; + # if present here it would still be a test failure if acted on. + ], + "merged_pr_count": 1, + }, + ).start() + patch( + "mcp_server.merged_cleanup_reconcile.remove_local_worktree", + side_effect=fake_remove, + ).start() + patch( + "mcp_server._collect_branch_ownership_records", + side_effect=fake_collect, + ).start() + patch("mcp_server._probe_remote_branch", side_effect=fake_probe).start() + + res = gitea_reconcile_merged_cleanups( + dry_run=False, + execute_confirmed=True, + pr_number=848, + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + ) + self.assertTrue(res.get("performed") or res.get("executed")) + self.assertEqual(res.get("selection_mode"), "exact_pr") + self.assertEqual(res.get("selected_pr_number"), 848) + actions = res.get("actions") or [] + pr_numbers_touched = { + a.get("pr_number") for a in actions if a.get("pr_number") is not None + } + self.assertTrue(pr_numbers_touched.issubset({None, 848}) or not pr_numbers_touched) + removes = [a for a in actions if a.get("action") == "remove_local_worktree"] + deletes = [a for a in actions if a.get("action") == "delete_remote_branch"] + self.assertEqual(len(removes), 1) + self.assertEqual(remove_calls[0]["branch"], "fix/issue-844-exclude-epic-containers") + self.assertEqual(len(deletes), 1) + self.assertTrue(deletes[0].get("success")) + self.assertTrue(deletes[0].get("after_worktree_removal")) + self.assertEqual(len(delete_api_calls), 1) + self.assertEqual( + ownership_branches, ["fix/issue-844-exclude-epic-containers"] + ) + + def test_exact_pr_unknown_fails_closed_without_mutation(self): + from mcp_server import gitea_reconcile_merged_cleanups + + def fake_api(method, url, *args, **kwargs): + if method == "GET" and "/pulls/99999" in url: + raise RuntimeError("HTTP 404 Not Found") + raise AssertionError(f"unexpected API call {method} {url}") + + self.mock_api.side_effect = fake_api + res = gitea_reconcile_merged_cleanups( + dry_run=True, + pr_number=99999, + remote="prgs", + ) + self.assertFalse(res.get("success")) + self.assertFalse(res.get("performed")) + self.assertEqual(res.get("blocker_kind"), "pr_unresolvable") + self.assertIn("99999", " ".join(res.get("reasons") or [])) + + def test_exact_pr_not_merged_fails_closed(self): + from mcp_server import gitea_reconcile_merged_cleanups + + def fake_api(method, url, *args, **kwargs): + if method == "GET" and url.rstrip("/").endswith("/pulls/900"): + return { + "number": 900, + "merged": False, + "merged_at": None, + "state": "open", + "head": {"ref": "feat/x", "sha": "a" * 40}, + } + raise AssertionError(f"unexpected {method} {url}") + + self.mock_api.side_effect = fake_api + res = gitea_reconcile_merged_cleanups( + dry_run=False, + execute_confirmed=True, + pr_number=900, + remote="prgs", + ) + self.assertFalse(res.get("success")) + self.assertFalse(res.get("performed")) + self.assertEqual(res.get("blocker_kind"), "pr_not_merged") + + def test_exact_pr_invalid_number_fails_closed(self): + from mcp_server import gitea_reconcile_merged_cleanups + + res = gitea_reconcile_merged_cleanups( + dry_run=True, + pr_number=0, + remote="prgs", + ) + self.assertFalse(res.get("success")) + self.assertEqual(res.get("blocker_kind"), "invalid_pr_number") + self.mock_api.assert_not_called() + + def test_batch_mode_still_works_without_pr_number(self): + """Unfiltered batch path remains backward compatible.""" + from mcp_server import gitea_reconcile_merged_cleanups + + self.mock_all.side_effect = lambda url, auth, limit=None: [] + self.mock_api.side_effect = lambda *a, **k: {} + patch( + "mcp_server.merged_cleanup_reconcile.build_reconciliation_report", + return_value={ + "entries": [], + "reviewer_scratch_entries": [], + "merged_pr_count": 0, + }, + ).start() + res = gitea_reconcile_merged_cleanups(dry_run=True, remote="prgs", limit=10) + self.assertTrue(res.get("success")) + self.assertEqual(res.get("selection_mode"), "batch") + self.assertIsNone(res.get("selected_pr_number")) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_dirty_orphan_worktree_recovery.py b/tests/test_dirty_orphan_worktree_recovery.py new file mode 100644 index 0000000..4e58b4a --- /dev/null +++ b/tests/test_dirty_orphan_worktree_recovery.py @@ -0,0 +1,483 @@ +"""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"}, + ) + # #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)) + 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) + key = "symlink-test" + jdir = tmp + 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) + + +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() diff --git a/tests/test_issue_855_expired_reviewer_reclaim.py b/tests/test_issue_855_expired_reviewer_reclaim.py new file mode 100644 index 0000000..288b7f7 --- /dev/null +++ b/tests/test_issue_855_expired_reviewer_reclaim.py @@ -0,0 +1,244 @@ +"""#855 AC4: an expired reviewer lease must not indefinitely protect an +already-merged branch when no live claimant exists. + +Two layers are covered: + +* ``branch_cleanup_guard.assess_expired_reviewer_lease_reclaim`` — the pure, + fail-closed reclaim decision. Every condition must be provably satisfied or + the lease keeps protecting the branch. +* ``gitea_mcp_server._collect_branch_ownership_records`` — the wiring that + supplies authoritative evidence (PR merged state, owner-process liveness, + competing ownership) to that decision, and flips an expired reviewer lease + to reclaimable only under the full policy. + +All inputs are fabricated; no real repository, lease, or credential is used. +""" + +import importlib +import unittest +from unittest.mock import patch + +import branch_cleanup_guard + +mcp_server = importlib.import_module("gitea_mcp_server") + +FAKE_AUTH = "token fake" +REMOTE = "prgs" +ORG = "Scaled-Tech-Consulting" +REPO = "Gitea-Tools" +HOST = "gitea.prgs.cc" +BRANCH = "feat/issue-638-webui-app-shell-phase1" +PR_NUMBER = 818 + + +class TestAssessExpiredReviewerLeaseReclaim(unittest.TestCase): + """Pure fail-closed reclaim decision (#855 AC4).""" + + def _call(self, **overrides): + base = dict( + role="reviewer", + status="expired", + pr_merged=True, + owner_pid_alive=False, + competing_active_claimant=False, + ) + base.update(overrides) + return branch_cleanup_guard.assess_expired_reviewer_lease_reclaim(**base) + + def test_full_policy_satisfied_allows_reclaim(self): + out = self._call() + self.assertTrue(out["reclaim_allowed"]) + self.assertEqual(out["reasons"], []) + self.assertEqual(out["decision"], "reclaim_expired_reviewer_lease") + + def test_stale_dead_process_reviewer_also_reclaimable(self): + out = self._call(status="stale_dead_process") + self.assertTrue(out["reclaim_allowed"]) + + def test_non_reviewer_role_never_reclaims(self): + for role in ("author", "merger", "controller", "reconciler", "unknown"): + with self.subTest(role=role): + out = self._call(role=role) + self.assertFalse(out["reclaim_allowed"]) + self.assertTrue(out["reasons"]) + self.assertEqual(out["decision"], "keep_protecting") + + def test_active_status_never_reclaims(self): + out = self._call(status="active") + self.assertFalse(out["reclaim_allowed"]) + + def test_pr_not_merged_blocks_reclaim(self): + out = self._call(pr_merged=False) + self.assertFalse(out["reclaim_allowed"]) + + def test_pr_merged_unknown_fails_closed(self): + out = self._call(pr_merged=None) + self.assertFalse(out["reclaim_allowed"]) + + def test_owner_process_alive_blocks_reclaim(self): + out = self._call(owner_pid_alive=True) + self.assertFalse(out["reclaim_allowed"]) + + def test_owner_liveness_unknown_fails_closed(self): + out = self._call(owner_pid_alive=None) + self.assertFalse(out["reclaim_allowed"]) + + def test_competing_active_claimant_blocks_reclaim(self): + out = self._call(competing_active_claimant=True) + self.assertFalse(out["reclaim_allowed"]) + + def test_competing_claimant_unknown_fails_closed(self): + out = self._call(competing_active_claimant=None) + self.assertFalse(out["reclaim_allowed"]) + + def test_reasons_never_leak_secrets(self): + out = self._call(role="author") + blob = " ".join(out["reasons"]).lower() + self.assertNotIn("token", blob) + self.assertNotIn("password", blob) + + +class _FakeLease(dict): + pass + + +class TestCollectorExpiredReviewerReclaimWiring(unittest.TestCase): + """`_collect_branch_ownership_records` supplies authoritative evidence and + flips an expired reviewer lease to reclaimable only under the full policy.""" + + def _run( + self, + *, + lease_role="reviewer", + lease_freshness="stale_dead_process", + owner_pid_alive=False, + pr_merged=True, + extra_leases=None, + worktree_on_branch=False, + ): + lease = _FakeLease( + role=lease_role, + work_kind="pr", + work_number=PR_NUMBER, + branch=BRANCH, + status="active", + owner_pid=999999, + remote=REMOTE, + org=ORG, + repo=REPO, + host=HOST, + freshness={ + "freshness": lease_freshness, + "owner_pid": 999999, + "owner_pid_alive": owner_pid_alive, + "expired_by_time": lease_freshness == "expired", + }, + ) + leases = [lease] + list(extra_leases or []) + + pr_payload = { + "number": PR_NUMBER, + "merged": pr_merged, + "merged_at": "2026-07-23T00:00:00Z" if pr_merged else None, + "head": {"ref": BRANCH}, + } + + def fake_api_request(method, url, *a, **k): + if method == "GET" and f"/pulls/{PR_NUMBER}" in url: + return pr_payload + raise AssertionError(f"unexpected api_request {method} {url}") + + wt_entries = [] + if worktree_on_branch: + wt_entries = [{"branch": BRANCH, "path": f"/x/branches/{BRANCH}"}] + + with patch.object( + mcp_server.lease_lifecycle, + "list_active_leases", + return_value={"leases": leases}, + ), patch.object( + mcp_server.control_plane_db, "get_db", return_value=object(), create=True + ), patch.object( + mcp_server.issue_lock_store, "iter_lock_files", return_value=[] + ), patch.object( + mcp_server.worktree_cleanup_audit, + "list_worktrees", + return_value=wt_entries, + ), patch.object( + mcp_server, "api_get_all", return_value=[] + ), patch.object( + mcp_server, "api_request", side_effect=fake_api_request + ): + return mcp_server._collect_branch_ownership_records( + remote=REMOTE, + host=HOST, + org=ORG, + repo=REPO, + branch=BRANCH, + pr_number=PR_NUMBER, + project_root="/x", + auth=FAKE_AUTH, + base_api="https://gitea.prgs.cc/api/v1/repos/x/y", + ) + + def _reviewer_records(self, bundle): + return [ + rec + for rec in bundle["records"] + if rec.get("category") + == branch_cleanup_guard.OWNERSHIP_CATEGORY_REVIEWER_LEASE + ] + + def test_merged_dead_uncontested_reviewer_lease_is_reclaimable(self): + bundle = self._run() + self.assertFalse(bundle["inventory_error"]) + recs = self._reviewer_records(bundle) + self.assertEqual(len(recs), 1) + self.assertTrue(recs[0]["reclaim_allowed"]) + # And the guard consequently does not block deletion on it. + ownership = branch_cleanup_guard.assess_active_branch_ownership( + remote=REMOTE, org=ORG, repo=REPO, branch=BRANCH, host=HOST, + records=bundle["records"], + ) + self.assertFalse(ownership["block"]) + + def test_unmerged_pr_keeps_reviewer_lease_protective(self): + bundle = self._run(pr_merged=False) + recs = self._reviewer_records(bundle) + self.assertEqual(len(recs), 1) + self.assertFalse(recs[0]["reclaim_allowed"]) + ownership = branch_cleanup_guard.assess_active_branch_ownership( + remote=REMOTE, org=ORG, repo=REPO, branch=BRANCH, host=HOST, + records=bundle["records"], + ) + self.assertTrue(ownership["block"]) + + def test_owner_process_alive_keeps_reviewer_lease_protective(self): + bundle = self._run(owner_pid_alive=True, lease_freshness="expired") + recs = self._reviewer_records(bundle) + self.assertFalse(recs[0]["reclaim_allowed"]) + + def test_competing_worktree_binding_keeps_reviewer_lease_protective(self): + bundle = self._run(worktree_on_branch=True) + recs = self._reviewer_records(bundle) + self.assertFalse(recs[0]["reclaim_allowed"]) + ownership = branch_cleanup_guard.assess_active_branch_ownership( + remote=REMOTE, org=ORG, repo=REPO, branch=BRANCH, host=HOST, + records=bundle["records"], + ) + self.assertTrue(ownership["block"]) + + def test_expired_author_lease_never_reclaimed_by_reviewer_policy(self): + bundle = self._run(lease_role="author") + author_recs = [ + rec + for rec in bundle["records"] + if rec.get("category") + == branch_cleanup_guard.OWNERSHIP_CATEGORY_AUTHOR_LEASE + ] + self.assertEqual(len(author_recs), 1) + self.assertFalse(author_recs[0]["reclaim_allowed"]) + + +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( diff --git a/tests/test_mcp_restart_paths.py b/tests/test_mcp_restart_paths.py new file mode 100644 index 0000000..cf0c5eb --- /dev/null +++ b/tests/test_mcp_restart_paths.py @@ -0,0 +1,146 @@ +"""Tests for the MCP restart-path inventory and guards (#657). + +Covers: +* the registry is well-formed and every path is classified; +* unknown restart attempts fail closed (AC "fail closed on unknown restart"); +* the previously-unguarded full-restart primitives stay guarded/absent + against the real source tree (AC "tests for at least one previously + unguarded path"); +* pkill of the daemon is still classified as contamination (#630, AC3); +* the inventory doc and module stay in lock-step. +""" + +import os +import tempfile +import unittest +from pathlib import Path + +import mcp_restart_paths as rp +import runtime_recovery_guard + +REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +DOC_PATH = os.path.join(REPO_ROOT, "docs", "mcp-restart-path-inventory.md") + + +class TestRegistryWellformed(unittest.TestCase): + def test_registry_is_wellformed(self): + # Must not raise. + rp.assert_registry_wellformed() + + def test_every_path_has_valid_classification(self): + for path in rp.iter_restart_paths(): + self.assertIn(path.classification, rp.VALID_CLASSIFICATIONS) + self.assertTrue(path.guard.strip(), path.path_id) + self.assertTrue(path.references, path.path_id) + self.assertTrue(path.locations, path.path_id) + + def test_ids_are_unique(self): + ids = [p.path_id for p in rp.iter_restart_paths()] + self.assertEqual(len(ids), len(set(ids))) + + def test_covers_every_classification(self): + present = {p.classification for p in rp.iter_restart_paths()} + self.assertEqual(present, set(rp.VALID_CLASSIFICATIONS)) + + +class TestUnknownAttemptFailsClosed(unittest.TestCase): + def test_unknown_path_raises(self): + with self.assertRaises(rp.UnknownRestartPathError): + rp.assert_restart_attempt_registered("totally_novel_restart_hack") + + def test_get_unknown_raises(self): + with self.assertRaises(rp.UnknownRestartPathError): + rp.get_restart_path("nope") + + def test_registered_attempt_returns_path(self): + path = rp.assert_restart_attempt_registered("manual_daemon_kill") + self.assertEqual(path.classification, rp.CLASS_FORBIDDEN) + + +class TestDaemonNeverSelfReplaces(unittest.TestCase): + """Previously-unguarded full-restart primitive: daemon self-replacement.""" + + def test_no_self_replacement_in_source(self): + # The live daemon modules must contain no os.execv/os.kill/os._exit + # self-restart call. Must not raise. + rp.assert_no_daemon_self_replacement(REPO_ROOT) + + def test_scanner_flags_injected_violation(self): + # Guard the guard: prove the scanner catches a real self-replace call. + with tempfile.TemporaryDirectory() as tmp: + bad = Path(tmp) / "gitea_mcp_server.py" + bad.write_text( + "import os\n" + "def restart():\n" + " os.execv('/usr/bin/python', ['python'])\n", + encoding="utf-8", + ) + found = rp.scan_daemon_self_replacement(tmp) + self.assertTrue(found) + with self.assertRaises(AssertionError): + rp.assert_no_daemon_self_replacement(tmp) + + def test_scanner_ignores_comment_and_docstring_mentions(self): + with tempfile.TemporaryDirectory() as tmp: + ok = Path(tmp) / "gitea_mcp_server.py" + ok.write_text( + "import os\n" + "# NOT os.execv() to re-point the interpreter here.\n" + '"""Never calls os._exit to restart."""\n' + "value = 1\n", + encoding="utf-8", + ) + self.assertEqual(rp.scan_daemon_self_replacement(tmp), []) + + +class TestLegacyAutoRestartHelperRemoved(unittest.TestCase): + """Previously-unguarded full-restart path: _trigger_mcp_auto_restart.""" + + def test_helper_absent_in_source(self): + # Must not raise: helper was removed in #685. + rp.assert_auto_restart_helper_absent(REPO_ROOT) + + def test_scanner_flags_reintroduced_helper(self): + with tempfile.TemporaryDirectory() as tmp: + bad = Path(tmp) / "mcp_server.py" + bad.write_text( + "def _trigger_mcp_auto_restart():\n return True\n", + encoding="utf-8", + ) + with self.assertRaises(AssertionError): + rp.assert_auto_restart_helper_absent(tmp) + + +class TestPkillStaysForbidden(unittest.TestCase): + """AC3: pkill of the daemon remains forbidden/contaminating (#630).""" + + def test_manual_daemon_kill_registered_as_forbidden(self): + path = rp.get_restart_path("manual_daemon_kill") + self.assertEqual(path.classification, rp.CLASS_FORBIDDEN) + + def test_pkill_classified_as_contamination(self): + assessment = runtime_recovery_guard.assess_recovery_command( + "pkill -f mcp_server.py" + ) + self.assertTrue(assessment["contaminated"]) + + def test_read_only_probe_not_contamination(self): + assessment = runtime_recovery_guard.assess_recovery_command( + "ps aux | grep mcp_server" + ) + self.assertFalse(assessment["contaminated"]) + + +class TestInventoryDocInSync(unittest.TestCase): + def test_doc_exists(self): + self.assertTrue(os.path.exists(DOC_PATH), DOC_PATH) + + def test_doc_mentions_every_path_id(self): + with open(DOC_PATH, encoding="utf-8") as handle: + doc = handle.read() + for path in rp.iter_restart_paths(): + self.assertIn(path.path_id, doc, f"doc missing {path.path_id}") + + +if __name__ == "__main__": + unittest.main() 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,