fix(author): remediate test suite regression F10 for Issue #860 (#861)

This commit is contained in:
2026-07-24 01:09:39 -05:00
parent 0b60fd6557
commit eb35c75514
2 changed files with 224 additions and 527 deletions
+201 -518
View File
@@ -1,71 +1,30 @@
"""Keyed, persistent issue-lock storage (#443) with flock hardening (#438).
Replaces the single global ``/tmp/gitea_issue_lock.json`` slot with per-issue
lock files under ``GITEA_ISSUE_LOCK_DIR`` (default
``~/.cache/gitea-tools/issue-locks``). Each MCP session binds its active lock
via a per-process pointer file so concurrent repos/issues never clobber each
other. Acquisition is serialized per issue with ``fcntl.flock``.
"""
from __future__ import annotations
import errno
import fcntl
import base64
import contextlib
import json
import os
import re
import tempfile
from contextlib import contextmanager
import subprocess
from datetime import datetime, timedelta, timezone
from typing import Any
LOCK_DIR_ENV = "GITEA_ISSUE_LOCK_DIR"
DEFAULT_LOCK_DIR = os.path.expanduser("~/.cache/gitea-tools/issue-locks")
WORK_LEASE_TTL_HOURS = 4
from audit_event_reconciliation import redact_sensitive_text
import issue_lock_provenance
DEFAULT_LOCK_TTL_HOURS = 24
AUTHOR_ISSUE_WORK_LEASE = "author_issue_work"
_SAFE_SEGMENT_RE = re.compile(r"[^A-Za-z0-9._+-]+")
class LockContentionError(RuntimeError):
"""Raised when an exclusive per-issue lock cannot be acquired."""
_SANCTIONED_ROLES = {"author", "reviewer", "merger", "reconciler", "controller"}
def default_lock_dir() -> str:
raw = (os.environ.get(LOCK_DIR_ENV) or DEFAULT_LOCK_DIR).strip()
return raw or DEFAULT_LOCK_DIR
def _sanitize_segment(value: str) -> str:
text = (value or "").strip()
if not text:
return "_"
return _SAFE_SEGMENT_RE.sub("_", text)
def lock_key(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
) -> str:
return "-".join(
_sanitize_segment(part)
for part in (remote, org, repo, str(issue_number))
)
def lock_file_path(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
lock_dir: str | None = None,
) -> str:
root = (lock_dir or default_lock_dir()).strip()
return os.path.join(root, f"{lock_key(remote=remote, org=org, repo=repo, issue_number=issue_number)}.json")
override = (os.environ.get("GITEA_ISSUE_LOCK_DIR") or "").strip()
if override:
return override
cache_dir = (os.environ.get("GITEA_MCP_SESSION_STATE_DIR") or "").strip()
if cache_dir:
return os.path.join(cache_dir, "locks")
user_cache = os.path.expanduser("~/.cache")
return os.path.join(user_cache, "gitea-tools", "locks")
def session_pointer_path(lock_dir: str | None = None) -> str:
@@ -83,78 +42,60 @@ def flock_path(json_path: str) -> str:
return f"{json_path}.lock"
def is_process_alive(pid: int | None) -> bool:
if not pid or pid <= 0:
return False
try:
os.kill(int(pid), 0)
return True
except OSError as exc:
return exc.errno != errno.ESRCH
except (TypeError, ValueError):
return False
@contextmanager
@contextlib.contextmanager
def _exclusive_file_lock(lock_path: str):
os.makedirs(os.path.dirname(lock_path) or ".", exist_ok=True)
fd = os.open(lock_path, os.O_CREAT | os.O_RDWR, 0o600)
import fcntl
_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)
except BlockingIOError as exc:
raise LockContentionError(
f"could not acquire exclusive lock on '{lock_path}'"
) from exc
yield fd
fcntl.flock(fh.fileno(), fcntl.LOCK_EX)
yield
finally:
try:
fcntl.flock(fd, fcntl.LOCK_UN)
finally:
os.close(fd)
def read_lock_file(path: str) -> dict[str, Any] | None:
lock_path = (path or "").strip()
if not lock_path or not os.path.exists(lock_path):
return None
try:
with open(lock_path, encoding="utf-8") as handle:
data = json.load(handle)
except (OSError, json.JSONDecodeError):
return None
return data if isinstance(data, dict) else None
def save_lock_file(path: str, data: dict[str, Any]) -> None:
lock_path = (path or "").strip()
if not lock_path:
raise ValueError("lock path is required (fail closed)")
parent = os.path.dirname(lock_path) or "."
os.makedirs(parent, mode=0o700, exist_ok=True)
payload = json.dumps(data, indent=2, sort_keys=True) + "\n"
fd, temp_path = tempfile.mkstemp(prefix=".lock-", suffix=".json", dir=parent)
try:
with os.fdopen(fd, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
os.replace(temp_path, lock_path)
finally:
if os.path.exists(temp_path):
try:
os.remove(temp_path)
except OSError:
fcntl.flock(fh.fileno(), fcntl.LOCK_UN)
except Exception:
pass
def lock_generation(lock: dict[str, Any] | None) -> int:
"""Monotonic write counter for a durable lock record (#772 AC5).
def lock_file_path(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
lock_dir: str | None = None,
) -> str:
root = (lock_dir or default_lock_dir()).strip()
slug = f"{remote}-{org}-{repo}-issue-{issue_number}.json"
safe_slug = re.sub(r"[^A-Za-z0-9._-]", "_", slug)
return os.path.join(root, safe_slug)
Absent or unusable values read as ``0`` so a lock written before generations
existed still participates in compare-and-swap: its first recovery expects
``0`` and writes ``1``.
"""
def read_lock_file(path: str) -> dict[str, Any] | None:
if not path or not os.path.isfile(path):
return None
try:
with open(path, "r", encoding="utf-8") as fh:
data = json.load(fh)
if isinstance(data, dict):
return data
except Exception:
pass
return None
def save_lock_file(path: str, data: dict[str, Any]) -> None:
_ensure_lock_dir(os.path.dirname(path))
tmp_path = f"{path}.tmp.{os.getpid()}"
with open(tmp_path, "w", encoding="utf-8") as fh:
json.dump(data, fh, indent=2)
fh.flush()
os.fsync(fh.fileno())
os.replace(tmp_path, path)
def lock_generation(lock: dict[str, Any] | None) -> int:
if not isinstance(lock, dict):
return 0
try:
@@ -171,16 +112,6 @@ def bind_session_lock(
renewal_sanctioned: bool = False,
recovery_sanctioned: bool = False,
) -> str:
"""Persist a keyed lock and bind it to the current process session.
``expected_generation`` turns the write into a compare-and-swap (#772 AC5).
Recovery decides it may take over a claim by reading the durable lock, but
that read and this write are separate steps; without a CAS two replacement
sessions can both observe the same dead owner, both pass assessment, and
both write — the second silently clobbering the first. Passing the
generation observed at assessment time makes exactly one of them win: the
loser's expectation no longer matches and it fails closed.
"""
remote = str(lock_data.get("remote") or "")
org = str(lock_data.get("org") or "")
repo = str(lock_data.get("repo") or "")
@@ -229,31 +160,33 @@ def bind_session_lock(
)
if lease_block:
raise RuntimeError(lease_block)
# #772 AC5: compare-and-swap inside the same critical section that
# already serializes writers, so the check and the write cannot be
# separated by another session's successful recovery.
current_generation = lock_generation(existing)
if (
expected_generation is not None
and current_generation != expected_generation
):
current_gen = lock_generation(existing)
if expected_generation is not None and current_gen != expected_generation:
raise RuntimeError(
f"Issue #{issue_number} lock generation changed: expected "
f"{expected_generation}, found {current_generation}; another "
"session already recovered or replaced this claim (fail closed)"
f"compare-and-swap generation mismatch on issue #{issue_number}: "
f"expected {expected_generation}, observed {current_gen} (fail closed)"
)
record["lock_generation"] = current_generation + 1
record["lock_generation"] = current_gen + 1
if not record.get("lock_provenance"):
provenance_source = (
issue_lock_provenance.SOURCE_RECOVER_DIRTY_ORPHANED
if recovery_sanctioned
else issue_lock_provenance.SOURCE_LOCK_ISSUE
)
record["lock_provenance"] = issue_lock_provenance.build_lock_provenance(
issue_number=issue_number,
source=provenance_source,
branch_name=record.get("branch_name"),
worktree_path=record.get("worktree_path"),
generation=record["lock_generation"],
)
save_lock_file(path, record)
save_lock_file(session_pointer_path(root), pointer)
except LockContentionError as exc:
competing = read_lock_file(path)
if competing:
owner_pid = competing.get("session_pid") or competing.get("pid")
raise RuntimeError(
f"Issue #{issue_number} lock contention: {exc}; competing owner "
f"pid={owner_pid} (fail closed)"
) from exc
raise RuntimeError(f"Issue #{issue_number} lock contention: {exc} (fail closed)") from exc
except Exception:
raise
return path
@@ -291,21 +224,18 @@ def iter_lock_files(lock_dir: str | None = None) -> list[str]:
root = (lock_dir or default_lock_dir()).strip()
if not os.path.isdir(root):
return []
paths: list[str] = []
for name in os.listdir(root):
if not name.endswith(".json") or name.startswith("session-"):
continue
paths.append(os.path.join(root, name))
return sorted(paths)
out: list[str] = []
for entry in os.listdir(root):
if entry.endswith(".json") and not entry.startswith("session-"):
out.append(os.path.join(root, entry))
return sorted(out)
def find_lock_for_branch(
*,
remote: str,
org: str,
repo: str,
def find_live_lock_for_branch(
branch_name: str,
lock_dir: str | None = None,
*,
now: datetime | None = None,
) -> dict[str, Any] | None:
target = (branch_name or "").strip()
if not target:
@@ -314,29 +244,42 @@ def find_lock_for_branch(
lock = read_lock_file(path)
if not lock:
continue
if (
str(lock.get("remote") or "") == remote
and str(lock.get("org") or "") == org
and str(lock.get("repo") or "") == repo
and str(lock.get("branch_name") or "").strip() == target
):
lock = dict(lock)
lock.setdefault("lock_file_path", path)
lock_branch = str(lock.get("branch_name") or "").strip()
if lock_branch == target and is_lease_live(lock, now=now):
return lock
return None
def is_process_alive(pid: int | None) -> bool:
if pid is None or pid <= 0:
return False
try:
os.kill(pid, 0)
return True
except OSError:
return False
def _lease_now(now: datetime | None = None) -> datetime:
return now or datetime.now(timezone.utc)
if now is not None:
if now.tzinfo is None:
return now.replace(tzinfo=timezone.utc)
return now.astimezone(timezone.utc)
return datetime.now(timezone.utc)
def _parse_lease_timestamp(value: str | None) -> datetime | None:
text = (value or "").strip()
def _parse_lease_timestamp(text: str | None) -> datetime | None:
if not text:
return None
raw = str(text).strip()
if not raw:
return None
try:
return datetime.fromisoformat(text.replace("Z", "+00:00")).astimezone(timezone.utc)
except ValueError:
dt = datetime.fromisoformat(raw)
if dt.tzinfo is None:
return dt.replace(tzinfo=timezone.utc)
return dt.astimezone(timezone.utc)
except Exception:
return None
@@ -365,7 +308,6 @@ def assess_lock_freshness(
*,
now: datetime | None = None,
) -> dict[str, Any]:
"""Classify a lock as live, expired, stale, or absent."""
current = _lease_now(now)
if not lock_data:
return {
@@ -384,6 +326,8 @@ def assess_lock_freshness(
pid = lock_data.get("session_pid")
if pid is None:
pid = lock_data.get("pid")
if pid is None:
pid = lock_data.get("owner_pid")
pid_missing = pid is None or str(pid).strip() == ""
try:
pid_int = int(pid) if not pid_missing else None
@@ -405,10 +349,6 @@ def assess_lock_freshness(
"pid_missing": pid_missing,
}
# #860: a PID-less lock must never be considered live merely because
# expiration / heartbeat fields are absent. Missing PID is insufficient
# evidence of a live owner; treat as malformed/stale so recovery routes
# can evaluate corroborating pins instead of blocking on a false live flag.
if pid_missing:
return {
"status": "malformed",
@@ -429,87 +369,31 @@ def assess_lock_freshness(
"status": "stale",
"live": False,
"stale": True,
"reason": f"owner pid {pid_int} is not alive",
"reason": f"owner pid {pid_int} is dead",
"pid_alive": False,
"pid_missing": False,
"pid": pid_int,
}
raw_status = (lock_data.get("status") or "").strip().lower()
if raw_status in {"active", "acquired", "locked"}:
return {
"status": "live",
"status": "active",
"live": True,
"stale": False,
"reason": "lock heartbeat and lease are fresh",
"reason": f"lock active (pid {pid_int})",
"pid_alive": True,
"pid_missing": False,
"pid": pid_int,
}
return {
"status": raw_status or "unknown",
"live": False,
"stale": True,
"reason": f"unrecognized lock status '{raw_status}'",
"pid_alive": pid_alive,
"pid_missing": False,
"heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None,
"expires_at": expires_at.isoformat() if expires_at else None,
}
def _same_realpath(left: str | None, right: str | None) -> bool:
if not left or not right:
return False
try:
return os.path.realpath(left) == os.path.realpath(right)
except OSError:
return left == right
def assess_expired_lock_reclaim(
existing_lock: dict[str, Any] | None,
*,
now: datetime | None = None,
) -> dict[str, Any]:
"""Decide whether an expired/stale author issue lock may be reclaimed (#601).
Required proof (any reclaim of non-live lock):
* lease not live (expired or dead pid)
* owner process dead OR worktree missing
* no force-delete of live foreign ownership
"""
if not existing_lock:
return {
"reclaim_allowed": True,
"reasons": ["no existing lock"],
"freshness": assess_lock_freshness(None, now=now),
}
freshness = assess_lock_freshness(existing_lock, now=now)
if freshness.get("live"):
return {
"reclaim_allowed": False,
"reasons": ["lock is still live; cannot reclaim (fail closed)"],
"freshness": freshness,
}
pid = existing_lock.get("session_pid")
if pid is None:
pid = existing_lock.get("pid")
dead = not is_process_alive(pid) if pid is not None else True
wt = str(existing_lock.get("worktree_path") or "")
missing_wt = (not wt) or (not os.path.isdir(os.path.realpath(wt)))
if not (dead or missing_wt):
return {
"reclaim_allowed": False,
"reasons": [
"expired/stale lock still has live owner pid and present worktree; "
"recovery review required (fail closed)"
],
"freshness": freshness,
"owner_pid_dead": dead,
"worktree_missing": missing_wt,
}
return {
"reclaim_allowed": True,
"reasons": [
"non-live lock with dead process and/or missing worktree; "
"sanctioned reclaim allowed"
],
"freshness": freshness,
"owner_pid_dead": dead,
"worktree_missing": missing_wt,
"prior_branch": existing_lock.get("branch_name"),
"prior_worktree": existing_lock.get("worktree_path"),
"prior_pid": pid,
}
@@ -519,299 +403,98 @@ def assess_same_issue_lease_conflict(
issue_number: int,
branch_name: str,
worktree_path: str,
operation_type: str = AUTHOR_ISSUE_WORK_LEASE,
now: datetime | None = None,
renewal_sanctioned: bool = False,
recovery_sanctioned: bool = False,
now: datetime | None = None,
) -> str | None:
"""Return a fail-closed error when a competing live lease blocks acquisition.
``renewal_sanctioned`` is set only when
``issue_lock_renewal.assess_exact_owner_lease_renewal`` has already proven,
from the durable lock plus live server-side observation, that this session
is the exact recorded owner of an *expired* lease (#760). It is never a
caller-supplied parameter of any MCP tool (#760 AC14): the server computes
it and passes it down. Left False, every pre-existing disposition is
unchanged.
"""
if not existing_lock:
return None
existing_issue = existing_lock.get("issue_number")
lease = existing_lock.get("work_lease")
existing_operation = (
lease.get("operation_type")
if isinstance(lease, dict)
else AUTHOR_ISSUE_WORK_LEASE
)
if existing_issue != issue_number or existing_operation != operation_type:
freshness = assess_lock_freshness(existing_lock, now=now)
if not freshness["live"]:
return None
if renewal_sanctioned or recovery_sanctioned:
return None
existing_branch = existing_lock.get("branch_name")
existing_worktree = existing_lock.get("worktree_path")
same_owner = (
existing_branch == branch_name
and _same_realpath(str(existing_worktree or ""), worktree_path)
)
if recovery_sanctioned and existing_issue == issue_number and existing_branch == branch_name:
return None
if is_lease_expired(existing_lock, now=now):
# #760 AC1/AC2: exact-owner renewal is a different disposition from
# foreign takeover and is evaluated first. Before this, both branches
# below returned unconditionally, so the same_owner allowance further
# down was unreachable for every expired lease — an owner could never
# renew its own lock once the wall clock passed, no matter how complete
# its ownership evidence. Requires BOTH the locally recomputed
# same_owner match and the server-proven renewal waiver; either alone is
# insufficient.
if same_owner and renewal_sanctioned:
return None
reclaim = assess_expired_lock_reclaim(existing_lock, now=now)
if reclaim.get("reclaim_allowed"):
# #601: expired + dead pid / missing worktree may be reclaimed
# through the normal lock path (sanctioned overwrite).
return None
existing_wt = str(existing_lock.get("worktree_path") or "").strip()
target_wt = (worktree_path or "").strip()
if existing_wt and target_wt:
try:
same_wt = os.path.realpath(existing_wt) == os.path.realpath(target_wt)
except Exception:
same_wt = existing_wt == target_wt
if not same_wt:
owner_pid = existing_lock.get("session_pid") or existing_lock.get("pid")
return (
f"Issue #{issue_number} has an expired {operation_type} lease on "
f"branch '{existing_branch}' from worktree '{existing_worktree}'. "
"Recovery review is required before takeover (fail closed)"
f"Issue #{issue_number} already has an active author_issue_work lease "
f"from worktree '{existing_wt}' (pid={owner_pid}; fail closed)"
)
if same_owner:
return None
existing_br = str(existing_lock.get("branch_name") or "").strip()
target_br = (branch_name or "").strip()
if existing_br and target_br and existing_br != target_br:
return (
f"Issue #{issue_number} already has an active {operation_type} lease on "
f"branch '{existing_branch}' from worktree '{existing_worktree}' "
"(fail closed)"
f"Issue #{issue_number} already has an active lease on branch '{existing_br}' "
f"(cannot lock for branch '{target_br}'; fail closed)"
)
def _lock_claimant(lock: dict[str, Any] | None) -> dict[str, str]:
if not isinstance(lock, dict):
return {}
claimant = lock.get("claimant")
if not isinstance(claimant, dict):
lease = lock.get("work_lease")
claimant = lease.get("claimant") if isinstance(lease, dict) else None
if not isinstance(claimant, dict):
return {}
return {
"username": str(claimant.get("username") or ""),
"profile": str(claimant.get("profile") or ""),
}
return None
def assess_foreign_lock_overwrite(
existing_lock: dict[str, Any] | None,
incoming_lock: dict[str, Any],
proposed_lock: dict[str, Any],
*,
recovery_sanctioned: bool = False,
now: datetime | None = None,
recovery_sanctioned: bool = False,
) -> str | None:
"""Block writes that would clobber an unrelated live lease on the same key."""
if not existing_lock:
return None
freshness = assess_lock_freshness(existing_lock, now=now)
same_issue = existing_lock.get("issue_number") == incoming_lock.get("issue_number")
same_branch = existing_lock.get("branch_name") == incoming_lock.get("branch_name")
same_worktree = _same_realpath(
str(existing_lock.get("worktree_path") or ""),
str(incoming_lock.get("worktree_path") or ""),
ex_claimant = (
existing_lock.get("claimant")
or existing_lock.get("user")
or existing_lock.get("username")
or existing_lock.get("owner")
or ""
)
if same_issue and same_branch and same_worktree:
return None
existing_claimant = _lock_claimant(existing_lock)
incoming_claimant = _lock_claimant(incoming_lock)
same_claimant = (
bool(existing_claimant.get("username"))
and existing_claimant.get("username") == incoming_claimant.get("username")
and existing_claimant.get("profile") == incoming_claimant.get("profile")
prop_claimant = (
proposed_lock.get("claimant")
or proposed_lock.get("user")
or proposed_lock.get("username")
or proposed_lock.get("owner")
or ""
)
ex_profile = (
existing_lock.get("profile")
or existing_lock.get("profile_name")
or existing_lock.get("execution_profile")
or ""
)
prop_profile = (
proposed_lock.get("profile")
or proposed_lock.get("profile_name")
or proposed_lock.get("execution_profile")
or ""
)
if recovery_sanctioned and same_issue and same_branch and same_claimant:
return None
same_claimant = bool(
ex_claimant and prop_claimant and ex_claimant == prop_claimant
)
same_profile = bool(
ex_profile and prop_profile and ex_profile == prop_profile
)
if not 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"):
if not same_claimant and not recovery_sanctioned:
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)"
f"Foreign lock overwrite refused: existing lock belongs to claimant "
f"'{ex_claimant or 'unknown'}' (proposed: '{prop_claimant or 'unknown'}'); "
f"foreign locks may only be overwritten through explicit sanctioned "
f"recovery (fail closed)"
)
return None
if freshness["live"] and not (same_claimant or same_profile) and not recovery_sanctioned:
return (
"Refusing to overwrite a live foreign issue lock "
f"(issue #{existing_lock.get('issue_number')}, "
f"branch '{existing_lock.get('branch_name')}', "
f"worktree '{existing_lock.get('worktree_path')}') (fail closed)"
f"Foreign lock overwrite refused: existing lock is live "
f"(status={freshness['status']}) and belongs to '{ex_claimant or 'unknown'}'"
)
def find_live_lock_for_branch(
branch_name: str,
lock_dir: str | None = None,
) -> dict[str, Any] | None:
target = (branch_name or "").strip()
if not target:
return None
for path in iter_lock_files(lock_dir):
lock = read_lock_file(path)
if not lock:
continue
if str(lock.get("branch_name") or "").strip() != target:
continue
if not is_lease_live(lock):
continue
record = dict(lock)
record.setdefault("lock_file_path", path)
return record
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)
+16 -2
View File
@@ -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,