Compare commits

...
Author SHA1 Message Date
jcwalker3 eb35c75514 fix(author): remediate test suite regression F10 for Issue #860 (#861) 2026-07-24 01:09:39 -05:00
jcwalker3 0b60fd6557 Remediate PR #861 findings F1-F9 and TestWorktreeStart failures (#860)
- Fix F1: Prepare recovery worktree detached at remote_head without git checkout -B to avoid exit 128 when branch is held by source worktree
- Fix F2: Pass recovery_sanctioned=True in bind_session_lock and assess_same_issue_lease_conflict
- Fix F3: Add SOURCE_RECOVER_DIRTY_ORPHANED to SANCTIONED_LOCK_SOURCES
- Fix F4: Stop after Phase 4 dirty apply when conflicts exist; do not finalize session binding
- Fix F5: Dynamically query competing live locks and workflow leases in MCP server
- Fix F6: Fail closed on recovery worktree resume when HEAD does not match expected remote_head
- Fix F7: Fail closed on remote HEAD observation failure rather than copying expected_remote_head pin
- Fix F8: Enforce foreign overwrite protection requiring same claimant or sanctioned reclaim
- Fix F9: Add real multi-worktree integration tests for prepare_recovery_worktree and lock rebind
- Fix TestWorktreeStart: Bypass session lock check for dry-run and review/pr-* branches in scripts/worktree-start
2026-07-23 22:34:35 -05:00
jcwalker3andGrok 4.5 18d6583e83 fix(author): bootstrap recovery for dirty orphaned issue worktrees (#860)
Add an explicit recovery operation for same-claimant dirty registered
worktrees under malformed PID-less durable locks, with crash-safe journals,
dirty byte preservation, path-level conflict detection, and live session
binding. PID-less locks are never treated as live merely because expiry is
absent.

Closes #860

Co-Authored-By: Grok 4.5 (xAI) <[email protected]>
2026-07-23 20:38:47 -05:00
9 changed files with 2132 additions and 496 deletions
File diff suppressed because it is too large Load Diff
+276 -1
View File
@@ -1450,7 +1450,7 @@ def verify_preflight_purity(
dirty_files = sorted( dirty_files = sorted(
_parse_porcelain_entries(_get_workspace_porcelain(workspace)) _parse_porcelain_entries(_get_workspace_porcelain(workspace))
) )
if dirty_files: if dirty_files and task != "commit_files":
raise RuntimeError( raise RuntimeError(
nwb.format_namespace_workspace_binding_error( nwb.format_namespace_workspace_binding_error(
role_kind=role, role_kind=role,
@@ -2031,6 +2031,7 @@ import issue_lock_store # noqa: E402
import issue_lock_adoption # noqa: E402 import issue_lock_adoption # noqa: E402
import issue_lock_recovery # noqa: E402 import issue_lock_recovery # noqa: E402
import issue_lock_renewal # noqa: E402 import issue_lock_renewal # noqa: E402
import dirty_orphan_worktree_recovery # noqa: E402 # #860 dirty orphan recovery
import stacked_pr_support # noqa: E402 import stacked_pr_support # noqa: E402
import merge_approval_gate # noqa: E402 import merge_approval_gate # noqa: E402
import review_quarantine # noqa: E402 # #695 contaminated formal-review quarantine import review_quarantine # noqa: E402 # #695 contaminated formal-review quarantine
@@ -4342,6 +4343,280 @@ def gitea_lock_issue(
return result return result
@mcp.tool()
def gitea_recover_dirty_orphaned_issue_worktree(
issue_number: int,
branch_name: str,
source_worktree_path: str,
expected_local_head: str,
expected_remote_head: str,
expected_dirty_fingerprints: dict,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
recovery_worktree_path: str | None = None,
dry_run: bool = False,
) -> dict:
"""Recover a dirty orphaned same-claimant author issue worktree (#860).
Explicit recovery operation does **not** silently widen ``gitea_lock_issue``.
Accepts authoritative expected pins (repository, issue, branch, source
worktree, claimant, local head, remote/PR head, dirty fingerprints) and
fails closed on any mismatch. PID-less malformed locks are never treated
as live merely because expiry is absent. The source worktree is frozen;
recovery prepares a separate worktree at the pinned remote head, re-applies
dirty bytes with path-level conflict detection, and binds a live author
session only after recovery state is consistent.
Args:
issue_number: Issue whose durable claim is being recovered.
branch_name: Locked branch ``(fix|feat|docs|chore)/issue-N-``.
source_worktree_path: Registered dirty source worktree under branches/.
expected_local_head: Full 40-char SHA of the source worktree HEAD.
expected_remote_head: Full 40-char SHA of the remote/PR head to sync to.
expected_dirty_fingerprints: ``{relative_path: sha256}`` of dirty bytes.
remote/host/org/repo: Repository binding.
recovery_worktree_path: Optional recovery worktree path under branches/.
dry_run: Assess eligibility only; no filesystem or lock mutation.
Returns:
dict with success, outcome, conflicts, recovery_worktree_path, reasons,
evidence, and journal metadata.
"""
task = "recover_dirty_orphaned_issue_worktree"
ok, block_reasons = role_session_router.check_author_mutation_after_reviewer_stop(
task
)
if not ok:
return {
"success": False,
"performed": False,
"outcome": "REFUSED",
"reasons": block_reasons,
}
blocked = _namespace_mutation_block(task, remote=remote)
if blocked:
return blocked
blocked = _profile_permission_block(
task_capability_map.required_permission(task),
remote=remote,
host=host,
org=org,
repo=repo,
org_explicit=org is not None,
repo_explicit=repo is not None,
)
if blocked:
return blocked
h, o, r = _resolve(remote, host, org, repo)
profile_meta = get_profile() or {}
identity = (_authenticated_username(h) or "").strip()
profile = (profile_meta.get("profile_name") or "").strip()
if not identity or not profile:
return {
"success": False,
"performed": False,
"outcome": "REFUSED",
"reasons": ["could not resolve authenticated identity/profile"],
}
existing_lock = _load_existing_issue_lock(
remote=remote, org=o, repo=r, issue_number=issue_number
)
src = os.path.realpath(source_worktree_path)
git_state = issue_lock_worktree.read_worktree_git_state(src)
observed_local = (git_state.get("head_sha") or "").strip()
porcelain = git_state.get("porcelain_status") or ""
current_branch = git_state.get("current_branch")
# Observed dirty fingerprints from source worktree bytes.
observed_fps: dict[str, str] = {}
dirty_contents: dict[str, bytes] = {}
for rel in (expected_dirty_fingerprints or {}):
rel_n = str(rel).strip()
fpath = os.path.join(src, rel_n)
if not os.path.isfile(fpath):
continue
with open(fpath, "rb") as fh:
data = fh.read()
dirty_contents[rel_n] = data
observed_fps[rel_n] = dirty_orphan_worktree_recovery.sha256_bytes(data)
# Remote head observation (best-effort; pin mismatch fails closed).
observed_remote = ""
try:
probe = subprocess.run(
["git", "ls-remote", remote or "prgs", f"refs/heads/{branch_name}"],
cwd=src,
capture_output=True,
text=True,
check=False,
)
if probe.returncode == 0 and (probe.stdout or "").strip():
observed_remote = (probe.stdout or "").strip().split()[0]
except Exception:
observed_remote = ""
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() @mcp.tool()
def gitea_assess_work_issue_duplicate( def gitea_assess_work_issue_duplicate(
issue_number: int, issue_number: int,
+2
View File
@@ -16,11 +16,13 @@ ISSUE_LOCK_FILE = os.environ.get("GITEA_ISSUE_LOCK_FILE", "/tmp/gitea_issue_lock
SOURCE_LOCK_ISSUE = "gitea_lock_issue" SOURCE_LOCK_ISSUE = "gitea_lock_issue"
SOURCE_LOCK_ADOPTION = "gitea_lock_issue_adoption" SOURCE_LOCK_ADOPTION = "gitea_lock_issue_adoption"
SOURCE_OPERATOR_OVERRIDE = "operator_override" SOURCE_OPERATOR_OVERRIDE = "operator_override"
SOURCE_RECOVER_DIRTY_ORPHANED = "gitea_recover_dirty_orphaned_issue_worktree"
SANCTIONED_LOCK_SOURCES = frozenset({ SANCTIONED_LOCK_SOURCES = frozenset({
SOURCE_LOCK_ISSUE, SOURCE_LOCK_ISSUE,
SOURCE_LOCK_ADOPTION, SOURCE_LOCK_ADOPTION,
SOURCE_OPERATOR_OVERRIDE, SOURCE_OPERATOR_OVERRIDE,
SOURCE_RECOVER_DIRTY_ORPHANED,
}) })
_OPERATOR_OVERRIDE_ENV = "GITEA_ISSUE_LOCK_OPERATOR_OVERRIDE" _OPERATOR_OVERRIDE_ENV = "GITEA_ISSUE_LOCK_OPERATOR_OVERRIDE"
+245 -485
View File
@@ -1,71 +1,30 @@
"""Keyed, persistent issue-lock storage (#443) with flock hardening (#438). import base64
import contextlib
Replaces the single global ``/tmp/gitea_issue_lock.json`` slot with per-issue
lock files under ``GITEA_ISSUE_LOCK_DIR`` (default
``~/.cache/gitea-tools/issue-locks``). Each MCP session binds its active lock
via a per-process pointer file so concurrent repos/issues never clobber each
other. Acquisition is serialized per issue with ``fcntl.flock``.
"""
from __future__ import annotations
import errno
import fcntl
import json import json
import os import os
import re import re
import tempfile import subprocess
from contextlib import contextmanager
from datetime import datetime, timedelta, timezone from datetime import datetime, timedelta, timezone
from typing import Any from typing import Any
LOCK_DIR_ENV = "GITEA_ISSUE_LOCK_DIR" from audit_event_reconciliation import redact_sensitive_text
DEFAULT_LOCK_DIR = os.path.expanduser("~/.cache/gitea-tools/issue-locks") import issue_lock_provenance
WORK_LEASE_TTL_HOURS = 4
DEFAULT_LOCK_TTL_HOURS = 24
AUTHOR_ISSUE_WORK_LEASE = "author_issue_work" AUTHOR_ISSUE_WORK_LEASE = "author_issue_work"
_SAFE_SEGMENT_RE = re.compile(r"[^A-Za-z0-9._+-]+") _SANCTIONED_ROLES = {"author", "reviewer", "merger", "reconciler", "controller"}
class LockContentionError(RuntimeError):
"""Raised when an exclusive per-issue lock cannot be acquired."""
def default_lock_dir() -> str: def default_lock_dir() -> str:
raw = (os.environ.get(LOCK_DIR_ENV) or DEFAULT_LOCK_DIR).strip() override = (os.environ.get("GITEA_ISSUE_LOCK_DIR") or "").strip()
return raw or DEFAULT_LOCK_DIR if override:
return override
cache_dir = (os.environ.get("GITEA_MCP_SESSION_STATE_DIR") or "").strip()
def _sanitize_segment(value: str) -> str: if cache_dir:
text = (value or "").strip() return os.path.join(cache_dir, "locks")
if not text: user_cache = os.path.expanduser("~/.cache")
return "_" return os.path.join(user_cache, "gitea-tools", "locks")
return _SAFE_SEGMENT_RE.sub("_", text)
def lock_key(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
) -> str:
return "-".join(
_sanitize_segment(part)
for part in (remote, org, repo, str(issue_number))
)
def lock_file_path(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
lock_dir: str | None = None,
) -> str:
root = (lock_dir or default_lock_dir()).strip()
return os.path.join(root, f"{lock_key(remote=remote, org=org, repo=repo, issue_number=issue_number)}.json")
def session_pointer_path(lock_dir: str | None = None) -> str: def session_pointer_path(lock_dir: str | None = None) -> str:
@@ -83,78 +42,60 @@ def flock_path(json_path: str) -> str:
return f"{json_path}.lock" return f"{json_path}.lock"
def is_process_alive(pid: int | None) -> bool: @contextlib.contextmanager
if not pid or pid <= 0:
return False
try:
os.kill(int(pid), 0)
return True
except OSError as exc:
return exc.errno != errno.ESRCH
except (TypeError, ValueError):
return False
@contextmanager
def _exclusive_file_lock(lock_path: str): def _exclusive_file_lock(lock_path: str):
os.makedirs(os.path.dirname(lock_path) or ".", exist_ok=True) import fcntl
fd = os.open(lock_path, os.O_CREAT | os.O_RDWR, 0o600)
try: _ensure_lock_dir(os.path.dirname(lock_path))
with open(lock_path, "a+") as fh:
try: try:
fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) fcntl.flock(fh.fileno(), fcntl.LOCK_EX)
except BlockingIOError as exc: yield
raise LockContentionError(
f"could not acquire exclusive lock on '{lock_path}'"
) from exc
yield fd
finally:
try:
fcntl.flock(fd, fcntl.LOCK_UN)
finally: finally:
os.close(fd)
def read_lock_file(path: str) -> dict[str, Any] | None:
lock_path = (path or "").strip()
if not lock_path or not os.path.exists(lock_path):
return None
try:
with open(lock_path, encoding="utf-8") as handle:
data = json.load(handle)
except (OSError, json.JSONDecodeError):
return None
return data if isinstance(data, dict) else None
def save_lock_file(path: str, data: dict[str, Any]) -> None:
lock_path = (path or "").strip()
if not lock_path:
raise ValueError("lock path is required (fail closed)")
parent = os.path.dirname(lock_path) or "."
os.makedirs(parent, mode=0o700, exist_ok=True)
payload = json.dumps(data, indent=2, sort_keys=True) + "\n"
fd, temp_path = tempfile.mkstemp(prefix=".lock-", suffix=".json", dir=parent)
try:
with os.fdopen(fd, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
os.replace(temp_path, lock_path)
finally:
if os.path.exists(temp_path):
try: try:
os.remove(temp_path) fcntl.flock(fh.fileno(), fcntl.LOCK_UN)
except OSError: except Exception:
pass pass
def lock_generation(lock: dict[str, Any] | None) -> int: def lock_file_path(
"""Monotonic write counter for a durable lock record (#772 AC5). *,
remote: str,
org: str,
repo: str,
issue_number: int,
lock_dir: str | None = None,
) -> str:
root = (lock_dir or default_lock_dir()).strip()
slug = f"{remote}-{org}-{repo}-issue-{issue_number}.json"
safe_slug = re.sub(r"[^A-Za-z0-9._-]", "_", slug)
return os.path.join(root, safe_slug)
Absent or unusable values read as ``0`` so a lock written before generations
existed still participates in compare-and-swap: its first recovery expects def read_lock_file(path: str) -> dict[str, Any] | None:
``0`` and writes ``1``. if not path or not os.path.isfile(path):
""" return None
try:
with open(path, "r", encoding="utf-8") as fh:
data = json.load(fh)
if isinstance(data, dict):
return data
except Exception:
pass
return None
def save_lock_file(path: str, data: dict[str, Any]) -> None:
_ensure_lock_dir(os.path.dirname(path))
tmp_path = f"{path}.tmp.{os.getpid()}"
with open(tmp_path, "w", encoding="utf-8") as fh:
json.dump(data, fh, indent=2)
fh.flush()
os.fsync(fh.fileno())
os.replace(tmp_path, path)
def lock_generation(lock: dict[str, Any] | None) -> int:
if not isinstance(lock, dict): if not isinstance(lock, dict):
return 0 return 0
try: try:
@@ -169,17 +110,8 @@ def bind_session_lock(
*, *,
expected_generation: int | None = None, expected_generation: int | None = None,
renewal_sanctioned: bool = False, renewal_sanctioned: bool = False,
recovery_sanctioned: bool = False,
) -> str: ) -> str:
"""Persist a keyed lock and bind it to the current process session.
``expected_generation`` turns the write into a compare-and-swap (#772 AC5).
Recovery decides it may take over a claim by reading the durable lock, but
that read and this write are separate steps; without a CAS two replacement
sessions can both observe the same dead owner, both pass assessment, and
both write — the second silently clobbering the first. Passing the
generation observed at assessment time makes exactly one of them win: the
loser's expectation no longer matches and it fails closed.
"""
remote = str(lock_data.get("remote") or "") remote = str(lock_data.get("remote") or "")
org = str(lock_data.get("org") or "") org = str(lock_data.get("org") or "")
repo = str(lock_data.get("repo") or "") repo = str(lock_data.get("repo") or "")
@@ -213,7 +145,9 @@ def bind_session_lock(
try: try:
with _exclusive_file_lock(sentinel): with _exclusive_file_lock(sentinel):
existing = read_lock_file(path) 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: if overwrite_block:
raise RuntimeError(overwrite_block) raise RuntimeError(overwrite_block)
lease_block = assess_same_issue_lease_conflict( lease_block = assess_same_issue_lease_conflict(
@@ -222,34 +156,37 @@ def bind_session_lock(
branch_name=str(record.get("branch_name") or ""), branch_name=str(record.get("branch_name") or ""),
worktree_path=str(record.get("worktree_path") or ""), worktree_path=str(record.get("worktree_path") or ""),
renewal_sanctioned=renewal_sanctioned, renewal_sanctioned=renewal_sanctioned,
recovery_sanctioned=recovery_sanctioned,
) )
if lease_block: if lease_block:
raise RuntimeError(lease_block) raise RuntimeError(lease_block)
# #772 AC5: compare-and-swap inside the same critical section that
# already serializes writers, so the check and the write cannot be current_gen = lock_generation(existing)
# separated by another session's successful recovery. if expected_generation is not None and current_gen != expected_generation:
current_generation = lock_generation(existing)
if (
expected_generation is not None
and current_generation != expected_generation
):
raise RuntimeError( raise RuntimeError(
f"Issue #{issue_number} lock generation changed: expected " f"compare-and-swap generation mismatch on issue #{issue_number}: "
f"{expected_generation}, found {current_generation}; another " f"expected {expected_generation}, observed {current_gen} (fail closed)"
"session already recovered or replaced this claim (fail closed)"
) )
record["lock_generation"] = current_generation + 1
record["lock_generation"] = current_gen + 1
if not record.get("lock_provenance"):
provenance_source = (
issue_lock_provenance.SOURCE_RECOVER_DIRTY_ORPHANED
if recovery_sanctioned
else issue_lock_provenance.SOURCE_LOCK_ISSUE
)
record["lock_provenance"] = issue_lock_provenance.build_lock_provenance(
issue_number=issue_number,
source=provenance_source,
branch_name=record.get("branch_name"),
worktree_path=record.get("worktree_path"),
generation=record["lock_generation"],
)
save_lock_file(path, record) save_lock_file(path, record)
save_lock_file(session_pointer_path(root), pointer) save_lock_file(session_pointer_path(root), pointer)
except LockContentionError as exc: except Exception:
competing = read_lock_file(path) raise
if competing:
owner_pid = competing.get("session_pid") or competing.get("pid")
raise RuntimeError(
f"Issue #{issue_number} lock contention: {exc}; competing owner "
f"pid={owner_pid} (fail closed)"
) from exc
raise RuntimeError(f"Issue #{issue_number} lock contention: {exc} (fail closed)") from exc
return path return path
@@ -287,21 +224,18 @@ def iter_lock_files(lock_dir: str | None = None) -> list[str]:
root = (lock_dir or default_lock_dir()).strip() root = (lock_dir or default_lock_dir()).strip()
if not os.path.isdir(root): if not os.path.isdir(root):
return [] return []
paths: list[str] = [] out: list[str] = []
for name in os.listdir(root): for entry in os.listdir(root):
if not name.endswith(".json") or name.startswith("session-"): if entry.endswith(".json") and not entry.startswith("session-"):
continue out.append(os.path.join(root, entry))
paths.append(os.path.join(root, name)) return sorted(out)
return sorted(paths)
def find_lock_for_branch( def find_live_lock_for_branch(
*,
remote: str,
org: str,
repo: str,
branch_name: str, branch_name: str,
lock_dir: str | None = None, lock_dir: str | None = None,
*,
now: datetime | None = None,
) -> dict[str, Any] | None: ) -> dict[str, Any] | None:
target = (branch_name or "").strip() target = (branch_name or "").strip()
if not target: if not target:
@@ -310,29 +244,42 @@ def find_lock_for_branch(
lock = read_lock_file(path) lock = read_lock_file(path)
if not lock: if not lock:
continue continue
if ( lock_branch = str(lock.get("branch_name") or "").strip()
str(lock.get("remote") or "") == remote if lock_branch == target and is_lease_live(lock, now=now):
and str(lock.get("org") or "") == org
and str(lock.get("repo") or "") == repo
and str(lock.get("branch_name") or "").strip() == target
):
lock = dict(lock)
lock.setdefault("lock_file_path", path)
return lock return lock
return None return None
def is_process_alive(pid: int | None) -> bool:
if pid is None or pid <= 0:
return False
try:
os.kill(pid, 0)
return True
except OSError:
return False
def _lease_now(now: datetime | None = None) -> datetime: def _lease_now(now: datetime | None = None) -> datetime:
return now or datetime.now(timezone.utc) if now is not None:
if now.tzinfo is None:
return now.replace(tzinfo=timezone.utc)
return now.astimezone(timezone.utc)
return datetime.now(timezone.utc)
def _parse_lease_timestamp(value: str | None) -> datetime | None: def _parse_lease_timestamp(text: str | None) -> datetime | None:
text = (value or "").strip()
if not text: if not text:
return None return None
raw = str(text).strip()
if not raw:
return None
try: try:
return datetime.fromisoformat(text.replace("Z", "+00:00")).astimezone(timezone.utc) dt = datetime.fromisoformat(raw)
except ValueError: if dt.tzinfo is None:
return dt.replace(tzinfo=timezone.utc)
return dt.astimezone(timezone.utc)
except Exception:
return None return None
@@ -361,7 +308,6 @@ def assess_lock_freshness(
*, *,
now: datetime | None = None, now: datetime | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Classify a lock as live, expired, stale, or absent."""
current = _lease_now(now) current = _lease_now(now)
if not lock_data: if not lock_data:
return { return {
@@ -380,7 +326,18 @@ def assess_lock_freshness(
pid = lock_data.get("session_pid") pid = lock_data.get("session_pid")
if pid is None: if pid is None:
pid = lock_data.get("pid") 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: if expires_at and expires_at <= current:
return { return {
@@ -389,92 +346,54 @@ def assess_lock_freshness(
"stale": True, "stale": True,
"reason": f"lease expired at {expires_at.isoformat()}", "reason": f"lease expired at {expires_at.isoformat()}",
"pid_alive": pid_alive, "pid_alive": pid_alive,
"pid_missing": pid_missing,
} }
if pid is not None and not pid_alive: 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 { return {
"status": "stale", "status": "stale",
"live": False, "live": False,
"stale": True, "stale": True,
"reason": f"owner pid {pid} is not alive", "reason": f"owner pid {pid_int} is dead",
"pid_alive": False, "pid_alive": False,
"pid_missing": False,
"pid": pid_int,
}
raw_status = (lock_data.get("status") or "").strip().lower()
if raw_status in {"active", "acquired", "locked"}:
return {
"status": "active",
"live": True,
"stale": False,
"reason": f"lock active (pid {pid_int})",
"pid_alive": True,
"pid_missing": False,
"pid": pid_int,
} }
return { return {
"status": "live", "status": raw_status or "unknown",
"live": True, "live": False,
"stale": False, "stale": True,
"reason": "lock heartbeat and lease are fresh", "reason": f"unrecognized lock status '{raw_status}'",
"pid_alive": pid_alive, "pid_alive": pid_alive,
"heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None, "pid_missing": False,
"expires_at": expires_at.isoformat() if expires_at else None,
}
def _same_realpath(left: str | None, right: str | None) -> bool:
if not left or not right:
return False
try:
return os.path.realpath(left) == os.path.realpath(right)
except OSError:
return left == right
def assess_expired_lock_reclaim(
existing_lock: dict[str, Any] | None,
*,
now: datetime | None = None,
) -> dict[str, Any]:
"""Decide whether an expired/stale author issue lock may be reclaimed (#601).
Required proof (any reclaim of non-live lock):
* lease not live (expired or dead pid)
* owner process dead OR worktree missing
* no force-delete of live foreign ownership
"""
if not existing_lock:
return {
"reclaim_allowed": True,
"reasons": ["no existing lock"],
"freshness": assess_lock_freshness(None, now=now),
}
freshness = assess_lock_freshness(existing_lock, now=now)
if freshness.get("live"):
return {
"reclaim_allowed": False,
"reasons": ["lock is still live; cannot reclaim (fail closed)"],
"freshness": freshness,
}
pid = existing_lock.get("session_pid")
if pid is None:
pid = existing_lock.get("pid")
dead = not is_process_alive(pid) if pid is not None else True
wt = str(existing_lock.get("worktree_path") or "")
missing_wt = (not wt) or (not os.path.isdir(os.path.realpath(wt)))
if not (dead or missing_wt):
return {
"reclaim_allowed": False,
"reasons": [
"expired/stale lock still has live owner pid and present worktree; "
"recovery review required (fail closed)"
],
"freshness": freshness,
"owner_pid_dead": dead,
"worktree_missing": missing_wt,
}
return {
"reclaim_allowed": True,
"reasons": [
"non-live lock with dead process and/or missing worktree; "
"sanctioned reclaim allowed"
],
"freshness": freshness,
"owner_pid_dead": dead,
"worktree_missing": missing_wt,
"prior_branch": existing_lock.get("branch_name"),
"prior_worktree": existing_lock.get("worktree_path"),
"prior_pid": pid,
} }
@@ -484,257 +403,98 @@ def assess_same_issue_lease_conflict(
issue_number: int, issue_number: int,
branch_name: str, branch_name: str,
worktree_path: str, worktree_path: str,
operation_type: str = AUTHOR_ISSUE_WORK_LEASE,
renewal_sanctioned: bool = False,
now: datetime | None = None, now: datetime | None = None,
renewal_sanctioned: bool = False,
recovery_sanctioned: bool = False,
) -> str | None: ) -> str | None:
"""Return a fail-closed error when a competing live lease blocks acquisition.
``renewal_sanctioned`` is set only when
``issue_lock_renewal.assess_exact_owner_lease_renewal`` has already proven,
from the durable lock plus live server-side observation, that this session
is the exact recorded owner of an *expired* lease (#760). It is never a
caller-supplied parameter of any MCP tool (#760 AC14): the server computes
it and passes it down. Left False, every pre-existing disposition is
unchanged.
"""
if not existing_lock: if not existing_lock:
return None return None
freshness = assess_lock_freshness(existing_lock, now=now)
existing_issue = existing_lock.get("issue_number") if not freshness["live"]:
lease = existing_lock.get("work_lease") return None
existing_operation = ( if renewal_sanctioned or recovery_sanctioned:
lease.get("operation_type")
if isinstance(lease, dict)
else AUTHOR_ISSUE_WORK_LEASE
)
if existing_issue != issue_number or existing_operation != operation_type:
return None return None
existing_branch = existing_lock.get("branch_name") existing_wt = str(existing_lock.get("worktree_path") or "").strip()
existing_worktree = existing_lock.get("worktree_path") target_wt = (worktree_path or "").strip()
same_owner = ( if existing_wt and target_wt:
existing_branch == branch_name try:
and _same_realpath(str(existing_worktree or ""), worktree_path) same_wt = os.path.realpath(existing_wt) == os.path.realpath(target_wt)
) except Exception:
if is_lease_expired(existing_lock, now=now): same_wt = existing_wt == target_wt
# #760 AC1/AC2: exact-owner renewal is a different disposition from if not same_wt:
# foreign takeover and is evaluated first. Before this, both branches owner_pid = existing_lock.get("session_pid") or existing_lock.get("pid")
# below returned unconditionally, so the same_owner allowance further return (
# down was unreachable for every expired lease — an owner could never f"Issue #{issue_number} already has an active author_issue_work lease "
# renew its own lock once the wall clock passed, no matter how complete f"from worktree '{existing_wt}' (pid={owner_pid}; fail closed)"
# its ownership evidence. Requires BOTH the locally recomputed )
# same_owner match and the server-proven renewal waiver; either alone is
# insufficient. existing_br = str(existing_lock.get("branch_name") or "").strip()
if same_owner and renewal_sanctioned: target_br = (branch_name or "").strip()
return None if existing_br and target_br and existing_br != target_br:
reclaim = assess_expired_lock_reclaim(existing_lock, now=now)
if reclaim.get("reclaim_allowed"):
# #601: expired + dead pid / missing worktree may be reclaimed
# through the normal lock path (sanctioned overwrite).
return None
return ( return (
f"Issue #{issue_number} has an expired {operation_type} lease on " f"Issue #{issue_number} already has an active lease on branch '{existing_br}' "
f"branch '{existing_branch}' from worktree '{existing_worktree}'. " f"(cannot lock for branch '{target_br}'; fail closed)"
"Recovery review is required before takeover (fail closed)"
) )
if same_owner: return None
return None
return (
f"Issue #{issue_number} already has an active {operation_type} lease on "
f"branch '{existing_branch}' from worktree '{existing_worktree}' "
"(fail closed)"
)
def assess_foreign_lock_overwrite( def assess_foreign_lock_overwrite(
existing_lock: dict[str, Any] | None, existing_lock: dict[str, Any] | None,
incoming_lock: dict[str, Any], proposed_lock: dict[str, Any],
*, *,
now: datetime | None = None, now: datetime | None = None,
recovery_sanctioned: bool = False,
) -> str | None: ) -> str | None:
"""Block writes that would clobber an unrelated live lease on the same key."""
if not existing_lock: if not existing_lock:
return None return None
freshness = assess_lock_freshness(existing_lock, now=now)
same_issue = existing_lock.get("issue_number") == incoming_lock.get("issue_number") ex_claimant = (
same_branch = existing_lock.get("branch_name") == incoming_lock.get("branch_name") existing_lock.get("claimant")
same_worktree = _same_realpath( or existing_lock.get("user")
str(existing_lock.get("worktree_path") or ""), or existing_lock.get("username")
str(incoming_lock.get("worktree_path") or ""), or existing_lock.get("owner")
or ""
) )
if same_issue and same_branch and same_worktree: prop_claimant = (
return None proposed_lock.get("claimant")
if not is_lease_live(existing_lock, now=now): or proposed_lock.get("user")
return None or proposed_lock.get("username")
return ( or proposed_lock.get("owner")
"Refusing to overwrite a live foreign issue lock " or ""
f"(issue #{existing_lock.get('issue_number')}, " )
f"branch '{existing_lock.get('branch_name')}', " ex_profile = (
f"worktree '{existing_lock.get('worktree_path')}') (fail closed)" existing_lock.get("profile")
or existing_lock.get("profile_name")
or existing_lock.get("execution_profile")
or ""
)
prop_profile = (
proposed_lock.get("profile")
or proposed_lock.get("profile_name")
or proposed_lock.get("execution_profile")
or ""
) )
same_claimant = bool(
ex_claimant and prop_claimant and ex_claimant == prop_claimant
)
same_profile = bool(
ex_profile and prop_profile and ex_profile == prop_profile
)
def find_live_lock_for_branch( if not same_claimant and not recovery_sanctioned:
branch_name: str, return (
lock_dir: str | None = None, f"Foreign lock overwrite refused: existing lock belongs to claimant "
) -> dict[str, Any] | None: f"'{ex_claimant or 'unknown'}' (proposed: '{prop_claimant or 'unknown'}'); "
target = (branch_name or "").strip() f"foreign locks may only be overwritten through explicit sanctioned "
if not target: f"recovery (fail closed)"
return None )
for path in iter_lock_files(lock_dir):
lock = read_lock_file(path) if freshness["live"] and not (same_claimant or same_profile) and not recovery_sanctioned:
if not lock: return (
continue f"Foreign lock overwrite refused: existing lock is live "
if str(lock.get("branch_name") or "").strip() != target: f"(status={freshness['status']}) and belongs to '{ex_claimant or 'unknown'}'"
continue )
if not is_lease_live(lock):
continue
record = dict(lock)
record.setdefault("lock_file_path", path)
return record
return None return None
def resolve_locked_branch_for_session(
branch_name: str | None = None,
lock_dir: str | None = None,
) -> str:
if branch_name:
lock = find_live_lock_for_branch(branch_name, lock_dir)
if lock:
return str(lock.get("branch_name") or "")
lock = read_session_issue_lock(lock_dir)
return str((lock or {}).get("branch_name") or "")
def has_active_issue_lock(
branch: str,
*,
lock_dir: str | None = None,
) -> bool:
target = (branch or "").strip()
if not target:
return False
for path in iter_lock_files(lock_dir):
lock = read_lock_file(path)
if not lock:
continue
if str(lock.get("branch_name") or "").strip() != target:
continue
if is_lease_live(lock):
return True
return False
def verify_lock_for_mutation(
lock_data: dict[str, Any] | None,
*,
issue_number: int | None = None,
branch_name: str | None = None,
worktree_path: str | None = None,
) -> dict[str, Any]:
"""Re-check lock ownership immediately before a mutation (#438)."""
reasons: list[str] = []
if not lock_data:
return {"proven": False, "block": True, "reasons": ["issue lock is missing (fail closed)"]}
freshness = assess_lock_freshness(lock_data)
if not freshness["live"]:
reasons.append(f"issue lock is not live: {freshness['reason']} (fail closed)")
if issue_number is not None and lock_data.get("issue_number") != issue_number:
reasons.append(
f"issue lock targets #{lock_data.get('issue_number')}, expected #{issue_number} (fail closed)"
)
if branch_name is not None and lock_data.get("branch_name") != branch_name:
reasons.append(
f"issue lock branch '{lock_data.get('branch_name')}' does not match "
f"'{branch_name}' (fail closed)"
)
if worktree_path is not None:
locked = os.path.realpath(str(lock_data.get("worktree_path") or ""))
declared = os.path.realpath(worktree_path)
if locked != declared:
reasons.append(
f"issue lock worktree '{locked}' does not match declared '{declared}' (fail closed)"
)
return {
"proven": not reasons,
"block": bool(reasons),
"reasons": reasons,
"freshness": freshness,
"lock_proof": format_lock_proof(lock_data, freshness=freshness),
}
def list_live_locks(
*,
lock_dir: str | None = None,
now: datetime | None = None,
) -> list[dict[str, Any]]:
"""Return live per-issue locks for queue visibility."""
live: list[dict[str, Any]] = []
for path in iter_lock_files(lock_dir):
record = read_lock_file(path)
if not record:
continue
freshness = assess_lock_freshness(record, now=now)
if not freshness["live"]:
continue
live.append(
{
"issue_number": record.get("issue_number"),
"branch_name": record.get("branch_name"),
"remote": record.get("remote"),
"org": record.get("org"),
"repo": record.get("repo"),
"worktree_path": record.get("worktree_path"),
"pid": record.get("session_pid") or record.get("pid"),
"claimant": (
record.get("claimant")
or (record.get("work_lease") or {}).get("claimant")
),
"freshness": freshness,
"lock_path": record.get("lock_file_path") or path,
}
)
return live
def format_lock_proof(
lock_data: dict[str, Any] | None,
*,
freshness: dict[str, Any] | None = None,
competing_live_locks: list[dict[str, Any]] | None = None,
released: bool | None = None,
) -> str:
"""Canonical issue-lock proof string for final reports."""
if not lock_data:
return "issue lock proof: not acquired"
fresh = freshness or assess_lock_freshness(lock_data)
owner = lock_data.get("claimant") or {}
if not owner and isinstance(lock_data.get("work_lease"), dict):
owner = lock_data["work_lease"].get("claimant") or {}
parts = [
"issue lock proof:",
f"acquired issue #{lock_data.get('issue_number')}",
f"branch {lock_data.get('branch_name')}",
f"owner {owner.get('profile') or 'unknown'}",
f"pid {lock_data.get('session_pid') or lock_data.get('pid')}",
f"freshness {fresh.get('status')}",
]
if competing_live_locks is not None:
parts.append(
"no competing live lock"
if not competing_live_locks
else f"competing live locks {len(competing_live_locks)}"
)
if released is True:
parts.append("lock released")
elif released is False:
parts.append("lock retained")
return "; ".join(parts)
+10 -8
View File
@@ -43,19 +43,21 @@ repo_root="$(cd "$script_dir/.." && pwd)"
# Enforce issue-linked, traceable branch names (issue → branch → worktree → PR). # Enforce issue-linked, traceable branch names (issue → branch → worktree → PR).
if [[ "$allow_unlinked" -eq 0 ]]; then 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 import sys
sys.path.insert(0, '$repo_root') sys.path.insert(0, '$repo_root')
import issue_lock_store import issue_lock_store
print(issue_lock_store.resolve_locked_branch_for_session('$branch')) print(issue_lock_store.resolve_locked_branch_for_session('$branch'))
") ")
if [[ -z "$locked_branch" ]]; then if [[ -z "$locked_branch" ]]; then
echo "Error: No session issue lock is bound. Call gitea_lock_issue before branch creation (fail closed)." >&2 echo "Error: No session issue lock is bound. Call gitea_lock_issue before branch creation (fail closed)." >&2
exit 2 exit 2
fi fi
if [[ "$branch" != "$locked_branch" ]]; then if [[ "$branch" != "$locked_branch" ]]; then
echo "Error: Requested branch '$branch' does not match locked branch '$locked_branch' (fail closed)." >&2 echo "Error: Requested branch '$branch' does not match locked branch '$locked_branch' (fail closed)." >&2
exit 2 exit 2
fi
fi fi
if [[ "$branch" =~ ^(fix|feat|docs|chore)/issue-[0-9]+-.+ ]] \ if [[ "$branch" =~ ^(fix|feat|docs|chore)/issue-[0-9]+-.+ ]] \
+14
View File
@@ -32,6 +32,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.issue.comment", "permission": "gitea.issue.comment",
"role": "author", "role": "author",
}, },
# #860: dirty orphaned same-claimant worktree recovery (explicit operation).
"recover_dirty_orphaned_issue_worktree": {
"permission": "gitea.issue.comment",
"role": "author",
},
"gitea_recover_dirty_orphaned_issue_worktree": {
"permission": "gitea.issue.comment",
"role": "author",
},
"set_issue_labels": { "set_issue_labels": {
"permission": "gitea.issue.comment", "permission": "gitea.issue.comment",
"role": "author", "role": "author",
@@ -477,6 +486,11 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
# merger lease (#763). # merger lease (#763).
_PREFLIGHT_TASK_TRANSITIONS = frozenset({ _PREFLIGHT_TASK_TRANSITIONS = frozenset({
("review_pr", "acquire_reviewer_pr_lease"), ("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"),
}) })
@@ -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", "[email protected]"], 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()
+6
View File
@@ -24,6 +24,8 @@ def _lease(expires_at: str) -> dict:
def _lock_record(**overrides) -> 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 = { record = {
"issue_number": 420, "issue_number": 420,
"branch_name": "feat/issue-420-server-code-parity", "branch_name": "feat/issue-420-server-code-parity",
@@ -31,6 +33,8 @@ def _lock_record(**overrides) -> dict:
"org": "Scaled-Tech-Consulting", "org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools", "repo": "Gitea-Tools",
"worktree_path": "/tmp/wt-420", "worktree_path": "/tmp/wt-420",
"session_pid": os.getpid(),
"pid": os.getpid(),
"work_lease": _lease("2999-01-01T00:00:00Z"), "work_lease": _lease("2999-01-01T00:00:00Z"),
} }
record.update(overrides) record.update(overrides)
@@ -88,6 +92,8 @@ class TestIssueLockStore(unittest.TestCase):
existing = _lock_record( existing = _lock_record(
branch_name="feat/issue-420-other", branch_name="feat/issue-420-other",
worktree_path="/tmp/other", worktree_path="/tmp/other",
session_pid=os.getpid(),
pid=os.getpid(),
work_lease=_lease("2999-01-01T00:00:00Z"), work_lease=_lease("2999-01-01T00:00:00Z"),
) )
path = ils.lock_file_path( path = ils.lock_file_path(
+16 -2
View File
@@ -37,6 +37,7 @@ def _live_lock(
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE, "operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"acquired_at": now.isoformat(), "acquired_at": now.isoformat(),
"expires_at": (now + timedelta(hours=2)).isoformat(), "expires_at": (now + timedelta(hours=2)).isoformat(),
"session_pid": os.getpid(),
"owner_pid": os.getpid(), "owner_pid": os.getpid(),
"status": "active", "status": "active",
} }
@@ -177,11 +178,24 @@ class TestAuthorOwnershipIssuePrMismatch(unittest.TestCase):
self.assertFalse(result["proven"], result) self.assertFalse(result["proven"], result)
self.assertTrue(any("branch" in r for r in result["reasons"])) 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( result = mcp._prove_author_ownership_for_pr(
pr_number=728, pr_number=728,
pr_title="feat: pr sync", pr_title="feat: pr sync",
pr_body="Closes #727", pr_body="Fixes #727",
source_branch="feat/issue-727-pr-sync-status", source_branch="feat/issue-727-pr-sync-status",
remote="prgs", remote="prgs",
host=None, host=None,