Compare commits

..
Author SHA1 Message Date
jcwalker3 6a636c58e7 Merge branch 'master' into fix/issue-790-slice-a-heartbeat-policy 2026-07-23 19:52:27 -05:00
jcwalker3 d2fe0110a0 Merge branch 'master' into fix/issue-790-slice-a-heartbeat-policy 2026-07-23 19:13:01 -05:00
jcwalker3 f1e4809930 Merge branch 'master' into fix/issue-790-slice-a-heartbeat-policy 2026-07-23 12:31:44 -05:00
jcwalker3 dc1d0e045f Merge branch 'master' into fix/issue-790-slice-a-heartbeat-policy 2026-07-23 01:13:11 -05:00
jcwalker3 badc4e636b Merge branch 'master' into fix/issue-790-slice-a-heartbeat-policy 2026-07-23 00:06:18 -05:00
sysadminandClaude Opus 4.8 243f52dc79 feat(lease): make the author task heartbeat load-bearing (#790 Slice A)
Slice A of Issue #790, per the controller reassessment in comment 13958. Does
not close the issue: terminal retirement (Slice B) and the read-side generation
check plus the #760 renewal re-scope (Slice C) are deliberately not implemented.

The defect. `issue_lock_store.assess_lock_freshness` parsed `last_heartbeat_at`
and then never consulted it. Liveness was decided by an absolute four-hour
`expires_at` and by PID liveness, and the recorded PID is the long-lived MCP
daemon rather than the authoring task, so an abandoned claim stayed live for the
full four hours. A tree-wide search found the field written in exactly one place
and advanced by nothing. Issue #787 / PR #789 hit this; Issue #760 / PR #791 hit
it again, blocking reconciliation for over five hours after its work had landed.

A1 — central policy. New `lease_policy` declares every duration for every task
class in one place: author initial/sliding TTL 10 minutes, heartbeat cadence 2,
stale warning 5, missed-heartbeat grace 10, absolute cap 8 hours, recovery grace
10, terminal race-drain 2. It ships first so the first heartbeat and TTL
behavior to run reads from it (AC-N7). The duplicated four-hour literal is gone
from both `issue_lock_store` and `gitea_mcp_server`. Reviewer, merger, and
conflict-fix classes are declared but not rewired — Slice C moves those call
sites — and a test asserts the declaration still equals the constants #747 and
`pr_work_lease` own, so the two cannot drift apart unnoticed.

A2 — load-bearing freshness, with two deliberate asymmetries. An alive PID never
establishes freshness anywhere (AC-N2); it is recorded as evidence and no branch
returns live because of it. A dead PID still marks a lease stale, and that band
still precedes every heartbeat evaluation, so #753 dead-session recovery keys on
exactly the classification it always did. New bands `stale_missed_heartbeat` and
`stale_absolute_cap` are classified in `branch_cleanup_guard` rather than
falling through to unknown-status, and still block unless the ownership record
proves `reclaim_allowed is True`. A heartbeat lease carrying no heartbeat is
contradictory and fails closed. `assess_expired_lock_reclaim` accepts a lapsed
heartbeat as reclaim grounds for heartbeat-lifecycle leases only: under this
lifecycle the heartbeat is the liveness proof, and also requiring a dead PID
would reinstate the original defect.

A3/A4 — task-session identity and the writer. `mint_task_session_id` produces an
ownership key containing no process identifier, since the daemon PID is reused
by every task it serves and identifies none of them. `heartbeat_session_lock`
writes inside the existing per-issue flock under the #772 generation
compare-and-swap, verifying exact issue, branch, realpath-normalized worktree,
claimant username, claimant profile, and recorded session identifier. It cannot
acquire, take over, or revive: a lease past its grace is refused and must use
the reclaim path, so a session that stopped proving liveness cannot restore
ownership retroactively. New `gitea_heartbeat_issue_lock` gates on the same
authority as `lock_issue`, being strictly narrower.

A5 — legacy compatibility (AC-N8). The explicit `lifecycle_version` marker, never
a timestamp comparison, discriminates legacy from heartbeat leases: a legacy lock
has `last_heartbeat_at == created_at` forever precisely because nothing advanced
it, and a freshly minted heartbeat lease has them equal too, so the equality
carries no information in either direction. Legacy locks keep their recorded
absolute expiry and are never evaluated against the short grace, so deployment
cannot make an existing claim instantly reclaimable. They leave that state only
by terminal retirement (Slice B) or by `rebind_legacy_lock`, which re-verifies
the exact owner and mints a genuine identifier and first heartbeat while
preserving the original claim under `legacy_origin`. Rebinding a lapsed legacy
lease is refused; that belongs to #760 renewal or #601 reclaim.

A6 — native coverage. Review #499 proved assessor-level tests miss discard
points, so `tests/test_issue_790_heartbeat_mcp_path.py` drives the real tools
against a real git repository and a real durable lock: lock creation and
read-back, policy window, freshness, survival of `verify_lock_for_mutation`,
invariance of the duplicate-work and linked-open-PR gates, CAS rejection,
foreign-session and foreign-claimant refusal, alive-PID-only refusal, missed
heartbeat, legacy protection on deployment, and legacy rebinding.

Tests. New suites 55 passed. Lock and lease regression set (issue_lock_store,
lease_lifecycle, #753, #755, #760 x2, #768, #772, lock registration, worktree,
adoption, duplicate gate, branch cleanup guard, capability invariants, claim
heartbeat, worktrees) 383 passed with 98 subtests. Full suite 4295 passed, 11
failed, 6 skipped, 499 subtests passed, against a clean master baseline worktree
at 620ed6e9 that reports 11 failed and 4240 passed — the same eleven node IDs.
The 55-test delta is exactly the new suites; no new failures.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_011u6GKSJwwrrYjguPjs1aK5
2026-07-22 04:19:48 -04:00
16 changed files with 1962 additions and 4184 deletions
+13 -85
View File
@@ -163,7 +163,19 @@ _TERMINAL_OWNERSHIP_STATUSES = frozenset(
{"released", "abandoned", "done", "blocked", "terminal", "closed"}
)
_EXPIRED_STATUSES = frozenset({"expired"})
_STALE_STATUSES = frozenset({"stale", "stale_dead_process", "stale_missing_worktree"})
_STALE_STATUSES = frozenset(
{
"stale",
"stale_dead_process",
"stale_missing_worktree",
# #790 Slice A heartbeat-lifecycle bands. Listed here so they are
# *classified* rather than falling through to the unknown-status branch;
# they still block unless the ownership record proves
# ``reclaim_allowed is True``, so the O2 fail-closed rule is unchanged.
"stale_missed_heartbeat",
"stale_absolute_cap",
}
)
def _norm_str(value: Any) -> str:
@@ -525,90 +537,6 @@ 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,
File diff suppressed because it is too large Load Diff
+1
View File
@@ -100,6 +100,7 @@ that gates each call, not which tools exist.
- `gitea_get_profile`
- `gitea_get_runtime_context`
- `gitea_get_shell_health`
- `gitea_heartbeat_issue_lock`
- `gitea_heartbeat_reviewer_pr_lease`
- `gitea_inspect_workflow_lease`
- `gitea_issue_irrecoverable_provenance_authorization`
+133 -538
View File
@@ -440,32 +440,13 @@ def _session_author_lock_worktree() -> str | None:
Used to derive the author mutation workspace when no explicit
``worktree_path`` or env binding is provided. Never invents a path.
#864: a session pointer whose owner PID is dead and is not this process
must not force workspace binding for other issues rebind is required for
that issue, and a stale dead-owner pointer must not poison unrelated work.
"""
try:
lock = issue_lock_store.read_session_issue_lock() or {}
except Exception:
return None
path = (lock.get("worktree_path") or "").strip()
if not path:
return None
pid = lock.get("session_pid")
if pid is None:
pid = lock.get("pid")
try:
pid_i = int(pid) if pid is not None else None
except (TypeError, ValueError):
pid_i = None
if (
pid_i is not None
and pid_i != os.getpid()
and not issue_lock_store.is_process_alive(pid_i)
):
return None
return path
return path or None
def _resolve_preflight_workspace_path(worktree_path: str | None = None) -> str:
@@ -2039,6 +2020,7 @@ import allocator_dependencies # noqa: E402
import dependency_graph # noqa: E402 # #784 durable dependency edges
import control_plane_db # noqa: E402
import lease_lifecycle # noqa: E402
import lease_policy # noqa: E402
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
import incident_bridge # noqa: E402
import sentry_observability # noqa: E402 (#606 optional Sentry observability)
@@ -2050,7 +2032,6 @@ 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_same_claimant_session_rebind # noqa: E402 # #864
import stacked_pr_support # noqa: E402
import merge_approval_gate # noqa: E402
import review_quarantine # noqa: E402 # #695 contaminated formal-review quarantine
@@ -2282,7 +2263,6 @@ import canonical_comment_validator as ccv # noqa: E402
# GITEA_ISSUE_LOCK_DIR, bound to the current MCP session via a per-PID pointer.
# Legacy global path retained only for test/doc references — do not seed manually.
ISSUE_LOCK_FILE = "/tmp/gitea_issue_lock.json"
WORK_LEASE_TTL_HOURS = 4
AUTHOR_ISSUE_WORK_LEASE = "author_issue_work"
VALID_WORK_LEASE_OPERATIONS = frozenset({
AUTHOR_ISSUE_WORK_LEASE,
@@ -2582,7 +2562,12 @@ def _build_author_issue_work_lease(
host: str | None,
) -> dict:
created = _work_lease_now()
expires = created + timedelta(hours=WORK_LEASE_TTL_HOURS)
# #790 Slice A: the window comes from the central policy, not a literal here.
# It is also now a *sliding* window — the lease lives ``initial_ttl_minutes``
# past its last valid heartbeat rather than a fixed four hours past its
# creation, so an abandoned task stops holding the claim within one TTL.
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK)
expires = created + timedelta(minutes=policy.initial_ttl_minutes)
return {
"operation_type": AUTHOR_ISSUE_WORK_LEASE,
"issue_number": issue_number,
@@ -2593,6 +2578,15 @@ def _build_author_issue_work_lease(
"created_at": _work_lease_timestamp(created),
"expires_at": _work_lease_timestamp(expires),
"last_heartbeat_at": _work_lease_timestamp(created),
# #790 AC-N1: the ownership key for this task. Distinct from the recorded
# PID, which is the shared daemon and identifies no individual task.
"task_session_id": issue_lock_store.mint_task_session_id(
AUTHOR_ISSUE_WORK_LEASE
),
# #790 AC-N8: the explicit lifecycle marker. Its absence — never a
# timestamp comparison — is what makes a lock legacy.
"lifecycle_version": lease_policy.LIFECYCLE_HEARTBEAT_V1,
"heartbeat_count": 1,
}
@@ -4363,260 +4357,135 @@ def gitea_lock_issue(
@mcp.tool()
def gitea_rebind_dirty_same_claimant_author_session(
def gitea_heartbeat_issue_lock(
issue_number: int,
branch_name: str,
worktree_path: str,
old_pid: int,
expected_local_head: str,
expected_remote_head: str,
expected_dirty_paths: list[str],
expected_fingerprints: dict,
task_session_id: str | None = None,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
dry_run: bool = False,
authorize_reconciler_execute: bool = False,
worktree_path: str | None = None,
expected_generation: int | None = None,
) -> dict:
"""Rebind a dirty registered issue worktree to this session (#864).
"""Prove an owned author issue lease is still active (#790 Slice A).
Sanctioned only when every pin agrees: same claimant, dead old_pid matching
the durable lock, matching local/remote heads, exact dirty path set, and
per-path sha256 fingerprints. Preserves every tracked/untracked byte.
Does not sync remote, create recovery worktrees, clean, reset, or move heads.
The task-liveness signal the lifecycle was missing. Before this, an author
lease carried a fixed four-hour expiry that nothing could shorten, and the
only liveness evidence was the recorded PID the long-lived MCP daemon,
which stays alive across every task it serves and so proved nothing about
whether the authoring task still held the work.
Role gate:
* author must match the lock claimant identity/profile
* reconciler execute only when ``authorize_reconciler_execute=True``
* reviewer/merger always refuse
Each successful call slides the lease ``initial_ttl_minutes`` past *now*
from the central policy, so an actively heartbeating session is never
evicted while an abandoned one releases its claim within one TTL.
``gitea.issue.comment`` (author map entry) is required for mutation; dry_run
still assesses fully but writes nothing. Permission alone is never ownership
proof every pin is re-checked server-side.
What this tool cannot do, by construction:
* **Acquire.** It refuses when no durable lock exists.
* **Take over.** Exact issue, branch, realpath-normalized worktree,
claimant username, claimant profile, and recorded task-session identifier
must all match; a superseded session holding an older identifier is
refused.
* **Revive.** A lease already past its grace is not heartbeatable that
would let a session restore ownership it had stopped proving. It must use
the sanctioned reclaim path, which mints a new generation.
A lock predating the heartbeat lifecycle is rebound rather than heartbeated:
its exact owner is re-verified and a genuine task-session identifier and
first heartbeat are minted (#790 AC-N8). The rebind is decided server-side
from the durable lifecycle marker; there is no caller-facing switch.
Args:
issue_number: Tracking issue number on the durable lock.
branch_name: Exact locked branch name.
worktree_path: Registered dirty worktree path (must be under branches/).
old_pid: Dead owner PID recorded on the lock (must match session_pid/pid).
expected_local_head: Full local HEAD sha the caller observed.
expected_remote_head: Full remote-tracking HEAD sha the caller observed.
expected_dirty_paths: Exact set of dirty relative paths (tracked+untracked).
expected_fingerprints: Map of relative path -> sha256 hex of file bytes.
issue_number: The locked issue number.
branch_name: The branch recorded on the lock.
task_session_id: The identifier this session received when it acquired
or rebound the lock. It is a fencing token, not an ownership
assertion: it is compared against durable state and can only ever
cause a refusal, never grant anything. Omitted only when rebinding a
legacy lock, which has no identifier yet and mints one.
remote: Known instance 'dadeschools' or 'prgs'.
host/org/repo: Optional target overrides (validated against binding).
dry_run: When true, assess only (no lock/session writes).
authorize_reconciler_execute: Reconciler-only execute gate.
host: Override the Gitea host.
org: Override the owner/organization.
repo: Override the repository name.
worktree_path: Author worktree recorded on the lock.
expected_generation: Optional fencing value. The per-issue flock already
serializes the read and the write, so this is for a caller that
wants to pin the generation it last observed across calls; a moved
generation fails closed.
Returns:
dict with 'success', 'performed', the sliding 'expires_at',
'last_heartbeat_at', 'lock_generation', 'task_session_id', the applied
'policy', and post-write 'freshness'; on refusal 'success'/'performed'
False with 'reasons' naming exactly what did not match.
"""
role = _profile_role_kind(get_profile())
role_norm = (role or "").strip().lower()
# Permission: authors need comment; dry_run assess is reachable under read
# for diagnosis, but execute always needs comment. Reconciler execute also
# needs comment when authorized.
if dry_run:
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"dry_run": True,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
else:
blocked = _profile_permission_block(
task_capability_map.required_permission(
"rebind_dirty_same_claimant_author_session"
),
issue_number=issue_number,
remote=remote,
host=host,
org=org,
repo=repo,
org_explicit=org is not None,
repo_explicit=repo is not None,
)
if blocked:
return blocked
if role_norm in {"reviewer", "merger"}:
return {
"success": False,
"dry_run": bool(dry_run),
"outcome": dirty_same_claimant_session_rebind.REFUSED,
"reasons": [
f"role '{role_norm}' cannot rebind dirty same-claimant author "
"sessions (fail closed)"
],
}
if role_norm == "reconciler" and not authorize_reconciler_execute and not dry_run:
return {
"success": False,
"dry_run": False,
"outcome": dirty_same_claimant_session_rebind.REFUSED,
"reasons": [
"reconciler role requires authorize_reconciler_execute=True "
"to execute dirty same-claimant rebind (fail closed)"
],
}
blocked = _profile_permission_block(
task_capability_map.required_permission("heartbeat_issue_lock"),
issue_number=issue_number,
remote=remote,
host=host,
org=org,
repo=repo,
org_explicit=org is not None,
repo_explicit=repo is not None,
)
if blocked:
return blocked
resolved_worktree = issue_lock_worktree.resolve_author_worktree_path(
worktree_path, _canonical_local_git_root()
)
h, o, r = _resolve(remote, host, org, repo)
try:
identity = _authenticated_username(h)
except Exception:
identity = None
profile = get_profile()
profile_name = profile.get("profile_name")
claimant = _work_lease_claimant(h)
identity = claimant.get("username")
profile = claimant.get("profile")
existing = _load_existing_issue_lock(
remote=remote, org=o, repo=r, issue_number=issue_number
)
resolved_wt = os.path.realpath(os.path.abspath((worktree_path or "").strip()))
inv = dirty_same_claimant_session_rebind.collect_dirty_inventory(resolved_wt)
if not existing:
return {
"success": False,
"performed": False,
"issue_number": issue_number,
"reasons": [
f"no durable lock for issue #{issue_number}; heartbeat cannot "
"acquire a claim (fail closed)"
],
}
branch_res = subprocess.run(
["git", "-C", resolved_wt, "branch", "--show-current"],
capture_output=True,
text=True,
check=False,
)
current_branch = (branch_res.stdout or "").strip() or None
head_res = subprocess.run(
["git", "-C", resolved_wt, "rev-parse", "HEAD"],
capture_output=True,
text=True,
check=False,
)
local_head = (head_res.stdout or "").strip() if head_res.returncode == 0 else None
# Observe remote-tracking head without network when possible.
remote_head = None
for ref in (
f"refs/remotes/origin/{branch_name}",
f"origin/{branch_name}",
f"refs/remotes/{remote}/{branch_name}",
f"{remote}/{branch_name}",
):
rh = subprocess.run(
["git", "-C", resolved_wt, "rev-parse", "--verify", "--quiet", ref],
capture_output=True,
text=True,
check=False,
if issue_lock_store.is_legacy_lease(existing):
# AC-N8 exit route one: canonical exact-owner rebinding. The other exit
# is terminal retirement, which is Slice B.
outcome = issue_lock_store.rebind_legacy_lock(
remote=remote,
org=o,
repo=r,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_worktree,
identity=identity,
profile=profile,
expected_generation=expected_generation,
)
if rh.returncode == 0 and (rh.stdout or "").strip():
remote_head = (rh.stdout or "").strip()
break
if remote_head is None:
# Fall back to caller's pin only for observation absence — assessment
# still requires pin==observed, so missing observation fails closed.
remote_head = None
outcome["operation"] = "legacy_rebind"
return outcome
# Competing live locks (other issues / other worktrees).
competing_live = []
for entry in issue_lock_store.list_live_locks():
competing_live.append(entry)
# Session pointers that claim this issue lock.
competing_sessions = []
lock_dir = issue_lock_store.default_lock_dir()
lock_path = issue_lock_store.lock_file_path(
remote=remote, org=o, repo=r, issue_number=issue_number, lock_dir=lock_dir
)
try:
for name in os.listdir(lock_dir):
if not name.startswith("session-") or not name.endswith(".json"):
continue
ptr = issue_lock_store.read_lock_file(os.path.join(lock_dir, name))
if not ptr:
continue
ptr_lock = str(ptr.get("lock_file_path") or "").strip()
if not ptr_lock:
continue
try:
same = os.path.realpath(ptr_lock) == os.path.realpath(lock_path)
except OSError:
same = ptr_lock == lock_path
if not same:
continue
try:
sess_pid = int(str(name)[len("session-") : -len(".json")])
except ValueError:
sess_pid = ptr.get("pid")
competing_sessions.append(
{
"pid": sess_pid,
"lock_file_path": ptr_lock,
"live": issue_lock_store.is_process_alive(sess_pid),
}
)
except OSError:
pass
# Best-effort workflow-lease scan: any live lock file whose work_lease is a
# non-author workflow lease on this issue/branch counts as active.
workflow_lease_active = False
for path in issue_lock_store.iter_lock_files(lock_dir):
rec = issue_lock_store.read_lock_file(path)
if not rec:
continue
lease = rec.get("work_lease") if isinstance(rec.get("work_lease"), dict) else {}
op = str(lease.get("operation_type") or "")
if op and op != issue_lock_store.AUTHOR_ISSUE_WORK_LEASE:
if rec.get("issue_number") == issue_number or str(
rec.get("branch_name") or ""
) == branch_name:
if issue_lock_store.is_lease_live(rec):
workflow_lease_active = True
break
repo_root = _canonical_local_git_root()
# permission_allowed reflects profile gate only — never ownership proof.
permission_allowed = True
result = dirty_same_claimant_session_rebind.apply_dirty_same_claimant_session_rebind(
outcome = issue_lock_store.heartbeat_session_lock(
remote=remote,
org=o,
repo=r,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_wt,
claimant_identity=identity,
claimant_profile=profile_name,
old_pid=old_pid,
expected_local_head=expected_local_head,
expected_remote_head=expected_remote_head,
expected_dirty_paths=list(expected_dirty_paths or []),
expected_fingerprints=dict(expected_fingerprints or {}),
existing_lock=existing,
current_identity=identity,
current_profile=profile_name,
role_kind=role_norm or role,
current_pid=os.getpid(),
current_branch=current_branch,
local_head=local_head,
remote_head=remote_head,
dirty_inventory=inv,
competing_live_locks=competing_live,
competing_sessions=competing_sessions,
workflow_lease_active=workflow_lease_active,
authorize_reconciler_execute=bool(authorize_reconciler_execute),
permission_allowed=permission_allowed,
repo_root=repo_root,
dry_run=bool(dry_run),
lock_dir=lock_dir,
worktree_path=resolved_worktree,
identity=identity,
profile=profile,
task_session_id=str(task_session_id or ""),
expected_generation=expected_generation,
)
result["observed"] = {
"local_head": local_head,
"remote_head": remote_head,
"current_branch": current_branch,
"dirty_paths": inv.get("dirty_paths"),
"fingerprints": inv.get("fingerprints"),
"identity": identity,
"profile": profile_name,
"role_kind": role_norm,
}
return result
outcome["operation"] = "heartbeat"
return outcome
@mcp.tool()
@@ -11116,9 +10985,6 @@ 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}
@@ -11267,28 +11133,15 @@ def _collect_branch_ownership_records(
else:
status = freshness_status
reclaim_allowed = False
rec = _base_rec(
category=category,
status=status,
reclaim_allowed=reclaim_allowed,
role=role,
host=lease_host or host_n or host,
records.append(
_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
@@ -11359,44 +11212,6 @@ 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}
@@ -11459,7 +11274,6 @@ 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,
@@ -11470,11 +11284,7 @@ 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 (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.
limit: Max number of closed PRs to inspect.
remote: Known Gitea instance ('dadeschools' or 'prgs').
host: Override the Gitea host.
org: Override the owner/organization.
@@ -11509,120 +11319,11 @@ 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)
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)
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] = {}
@@ -11649,13 +11350,6 @@ 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:
@@ -11690,33 +11384,6 @@ 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
@@ -12113,7 +11780,6 @@ def gitea_audit_worktree_cleanup(
org: str | None = None,
repo: str | None = None,
ttl_hours: float = worktree_cleanup_audit.DEFAULT_TTL_HOURS,
merged_pr_limit: int = 200,
) -> dict:
"""Read-only: classify every session-owned worktree under ``branches/`` (#401).
@@ -12124,26 +11790,17 @@ def gitea_audit_worktree_cleanup(
the active issue-lock branch is read from the local lock file and treated
as active work. Deletes nothing and mutates no Gitea state.
Merged PRs are fetched as well, so an issue worktree can be linked to the
PR that owns its branch (#858). Such a worktree only becomes removable
when that owning PR is unambiguous and merged, the worktree head is
already contained in authoritative master, and nothing else protects it
no open or competing PR, lease, issue lock, live session, dirty file, or
protected/control checkout. Anything unproven keeps it classified as
active issue work.
Fails closed if the live open-PR list, the merged-PR list, or the
control-plane lease state cannot be read: without them removability
cannot be proven, so no candidates are returned.
Fails closed if the live open-PR list cannot be fetched: without it,
removability cannot be proven, so no candidates are returned.
Args:
remote: Known instance 'dadeschools' or 'prgs'.
host: Override the Gitea host.
org: Override the owner/organization.
repo: Override the repository name.
ttl_hours: Age (hours) after which a clean conflict-fix worktree
becomes stale-removable (default from GITEA_WORKTREE_TTL_HOURS).
merged_pr_limit: Max closed PRs scanned for merged-PR ownership.
ttl_hours: Age (hours) after which a clean issue/conflict-fix
worktree becomes stale-removable (default from
GITEA_WORKTREE_TTL_HOURS).
Returns:
dict with per-worktree classifications, counts, removable
@@ -12179,84 +11836,22 @@ def gitea_audit_worktree_cleanup(
if (pr.get("head") or {}).get("ref")
}
# #858: merged PRs are the ownership evidence that lets a landed issue
# worktree stop being reported as active work. Without them the audit can
# never agree with the PR-scoped reconciler, so treat a fetch failure the
# same way an open-PR fetch failure is treated: fail closed.
try:
closed_prs = api_get_all(
f"{repo_api_url(h, o, r)}/pulls?state=closed", auth, limit=merged_pr_limit
)
except Exception as exc:
return {
"success": False,
"performed": False,
"open_pr_state_verified": True,
"merged_pr_state_verified": False,
"reasons": [
"could not fetch merged PRs; worktree ownership unverified "
f"(fail closed): {_redact(str(exc))}"
],
}
merged_prs = [pr for pr in closed_prs if (pr.get("merged") or pr.get("merged_at"))]
pr_index = worktree_cleanup_audit.build_pr_index(list(open_prs) + merged_prs)
# #858: the auditor already accepted lease evidence but nothing ever
# supplied it, so every worktree looked unleased. Removability is now
# reachable for issue worktrees, so authoritative control-plane leases
# must be readable or the audit fails closed.
db, lease_errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"performed": False,
"open_pr_state_verified": True,
"merged_pr_state_verified": True,
"lease_state_verified": False,
"reasons": [
"could not read control-plane leases; worktree protection "
"unverified (fail closed)",
*lease_errs,
],
}
lease_result = lease_lifecycle.list_active_leases(
db, remote=remote, org=o, repo=r, include_non_active=False, limit=500
)
leased_issue_numbers: set[int] = set()
live_session_paths: set[str] = set()
for lease in lease_result.get("leases") or []:
if lease.get("work_kind") == "issue" and lease.get("work_number") is not None:
try:
leased_issue_numbers.add(int(lease["work_number"]))
except (TypeError, ValueError):
pass
if lease.get("worktree_path"):
live_session_paths.add(str(lease["worktree_path"]))
active_issue_branches: set[str] = set()
lock = merged_cleanup_reconcile.read_issue_lock(ISSUE_LOCK_FILE)
if lock and lock.get("branch_name"):
active_issue_branches.add(str(lock["branch_name"]).strip())
master_ref = f"{remote}/master" if remote in REMOTES else "origin/master"
report = worktree_cleanup_audit.audit_branches_directory(
_canonical_local_git_root(),
open_pr_branches=open_pr_branches,
active_issue_branches=active_issue_branches,
now=datetime.now(timezone.utc),
ttl_hours=ttl_hours,
pr_index=pr_index,
leased_issue_numbers=leased_issue_numbers,
live_session_paths=live_session_paths,
master_ref=master_ref,
)
return {
"success": True,
"performed": False,
"open_pr_state_verified": True,
"merged_pr_state_verified": True,
"lease_state_verified": True,
"master_ref": master_ref,
"task_mode": "work-issue",
**report,
}
-5
View File
@@ -16,16 +16,11 @@ 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"
# #864: dirty-preserving same-claimant author-session rebind (dead owner PID).
SOURCE_DIRTY_SAME_CLAIMANT_REBIND = (
"gitea_rebind_dirty_same_claimant_author_session"
)
SANCTIONED_LOCK_SOURCES = frozenset({
SOURCE_LOCK_ISSUE,
SOURCE_LOCK_ADOPTION,
SOURCE_OPERATOR_OVERRIDE,
SOURCE_DIRTY_SAME_CLAIMANT_REBIND,
})
_OPERATOR_OVERRIDE_ENV = "GITEA_ISSUE_LOCK_OPERATOR_OVERRIDE"
+553 -29
View File
@@ -15,15 +15,27 @@ import json
import os
import re
import tempfile
import uuid
from contextlib import contextmanager
from datetime import datetime, timedelta, timezone
from typing import Any
import lease_policy
LOCK_DIR_ENV = "GITEA_ISSUE_LOCK_DIR"
DEFAULT_LOCK_DIR = os.path.expanduser("~/.cache/gitea-tools/issue-locks")
WORK_LEASE_TTL_HOURS = 4
AUTHOR_ISSUE_WORK_LEASE = "author_issue_work"
# Freshness classifications. ``STATUS_STALE`` remains the dead-PID band that
# #753 recovery keys on; the two bands below are new in #790 Slice A and apply
# only to leases minted under the heartbeat lifecycle.
STATUS_LIVE = "live"
STATUS_EXPIRED = "expired"
STATUS_ABSENT = "absent"
STATUS_STALE = "stale"
STATUS_STALE_MISSED_HEARTBEAT = "stale_missed_heartbeat"
STATUS_STALE_ABSOLUTE_CAP = "stale_absolute_cap"
_SAFE_SEGMENT_RE = re.compile(r"[^A-Za-z0-9._+-]+")
@@ -253,6 +265,331 @@ def bind_session_lock(
return path
def _ownership_refusals(
lock: dict[str, Any],
*,
issue_number: int,
branch_name: str,
worktree_path: str,
identity: str | None,
profile: str | None,
) -> list[str]:
"""Exact-ownership mismatches between a durable lock and a live caller.
Shared by the heartbeat writer and the legacy rebind path so the two cannot
disagree about what "the same owner" means. Every field is compared against
durable state; nothing is taken on the caller's word beyond the identity the
server itself resolved.
"""
reasons: list[str] = []
if lock.get("issue_number") != issue_number:
reasons.append(
f"lock targets issue #{lock.get('issue_number')}, not #{issue_number}"
)
if str(lock.get("branch_name") or "") != str(branch_name or ""):
reasons.append(
f"lock branch '{lock.get('branch_name')}' does not match '{branch_name}'"
)
if not _same_realpath(str(lock.get("worktree_path") or ""), worktree_path):
reasons.append(
f"lock worktree '{lock.get('worktree_path')}' does not match "
f"'{worktree_path}'"
)
lease = lock.get("work_lease") if isinstance(lock, dict) else None
claimant = lease.get("claimant") if isinstance(lease, dict) else None
claimant = claimant if isinstance(claimant, dict) else {}
recorded_identity = str(claimant.get("username") or "").strip()
recorded_profile = str(claimant.get("profile") or "").strip()
if not recorded_identity or not recorded_profile:
reasons.append("lock does not record both a claimant username and profile")
if recorded_identity and recorded_identity != str(identity or "").strip():
reasons.append(
f"lock claimant '{recorded_identity}' does not match active identity "
f"'{str(identity or '').strip() or 'unknown'}'"
)
if recorded_profile and recorded_profile != str(profile or "").strip():
reasons.append(
f"lock profile '{recorded_profile}' does not match active profile "
f"'{str(profile or '').strip() or 'unknown'}'"
)
return reasons
def _refusal(reasons: list[str], **extra: Any) -> dict[str, Any]:
return {"success": False, "performed": False, "reasons": reasons, **extra}
def heartbeat_session_lock(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
branch_name: str,
worktree_path: str,
identity: str | None,
profile: str | None,
task_session_id: str,
expected_generation: int | None = None,
lock_dir: str | None = None,
now: datetime | None = None,
) -> dict[str, Any]:
"""Slide a heartbeat-lifecycle lease forward (#790 Slice A, A4).
The write happens inside the same per-issue ``flock`` that serializes
acquisition, and under the #772 generation compare-and-swap, so a heartbeat
can never race a concurrent reclaim: whichever lands first moves the
generation and the other fails closed.
Refuses — never revives — in every ambiguous case. A lease that has already
lapsed past its grace is *not* heartbeatable: allowing that would let a
session that stopped proving liveness restore ownership retroactively, which
is precisely the revival AC-N5 forbids. Such a session must go through the
sanctioned reclaim path, which mints a fresh generation.
"""
current = _lease_now(now)
root = _ensure_lock_dir(lock_dir)
path = lock_file_path(
remote=remote, org=org, repo=repo, issue_number=issue_number, lock_dir=root
)
declared_session = str(task_session_id or "").strip()
if not declared_session:
return _refusal(["no task_session_id supplied (fail closed)"])
sentinel = flock_path(path)
try:
with _exclusive_file_lock(sentinel):
lock = read_lock_file(path)
if not lock:
return _refusal([f"no durable lock for issue #{issue_number}"])
if is_legacy_lease(lock):
return _refusal(
[
"lock predates the heartbeat lifecycle; it must be rebound "
"by its exact owner before it can be heartbeated"
],
lifecycle=lease_lifecycle_version(lock),
legacy_lease=True,
)
reasons = _ownership_refusals(
lock,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=worktree_path,
identity=identity,
profile=profile,
)
recorded_session = lease_task_session_id(lock)
if not recorded_session:
reasons.append(
"lock declares the heartbeat lifecycle but records no "
"task_session_id (fail closed)"
)
elif recorded_session != declared_session:
# A superseded session holding an old identifier cannot heartbeat
# over the session that replaced it.
reasons.append(
"task_session_id does not match the session recorded on the lock"
)
if reasons:
return _refusal(reasons)
current_generation = lock_generation(lock)
if (
expected_generation is not None
and current_generation != expected_generation
):
return _refusal(
[
f"lock generation changed: expected {expected_generation}, "
f"found {current_generation}; another session reclaimed or "
"replaced this claim (fail closed)"
],
lock_generation=current_generation,
)
freshness = assess_lock_freshness(lock, now=current)
if not freshness.get("live"):
return _refusal(
[
f"lease is not live ({freshness.get('status')}): "
f"{freshness.get('reason')}; a lapsed lease must be "
"reclaimed, not heartbeated"
],
freshness=freshness,
)
policy = lease_policy.policy_for(lease_task_class(lock))
expires = current + timedelta(minutes=policy.initial_ttl_minutes)
record = dict(lock)
lease = dict(record.get("work_lease") or {})
prior_heartbeat = lease.get("last_heartbeat_at")
lease["last_heartbeat_at"] = _format_lease_timestamp(current)
lease["expires_at"] = _format_lease_timestamp(expires)
try:
lease["heartbeat_count"] = int(lease.get("heartbeat_count") or 0) + 1
except (TypeError, ValueError):
lease["heartbeat_count"] = 1
record["work_lease"] = lease
record["lock_generation"] = current_generation + 1
save_lock_file(path, record)
except LockContentionError as exc:
return _refusal([f"issue #{issue_number} lock contention: {exc} (fail closed)"])
return {
"success": True,
"performed": True,
"issue_number": issue_number,
"branch_name": branch_name,
"worktree_path": worktree_path,
"task_session_id": declared_session,
"lock_generation": record["lock_generation"],
"prior_generation": current_generation,
"prior_heartbeat_at": prior_heartbeat,
"last_heartbeat_at": lease["last_heartbeat_at"],
"expires_at": lease["expires_at"],
"heartbeat_count": lease["heartbeat_count"],
"lock_file_path": path,
"policy": lease_policy.describe(lease_task_class(record)),
"freshness": assess_lock_freshness(record, now=current),
}
def rebind_legacy_lock(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
branch_name: str,
worktree_path: str,
identity: str | None,
profile: str | None,
expected_generation: int | None = None,
lock_dir: str | None = None,
now: datetime | None = None,
) -> dict[str, Any]:
"""Move a legacy lock into the heartbeat lifecycle (#790 AC-N8).
One of the two sanctioned exits from the preserved-expiry legacy state; the
other is terminal retirement, which is Slice B. Only the exact recorded
owner may rebind, and only while the legacy lock is still live under its
original absolute expiry — an already-expired legacy lease belongs to the
#760 renewal path or #601 reclaim, and this must not become a second, weaker
way to revive one.
The rebind mints a genuine task-session identifier and a genuine first
heartbeat. It does not fabricate history: the original creation and expiry
are preserved under ``legacy_origin`` for audit, and the new lifecycle's
absolute cap runs from the rebind, not from the legacy claim.
"""
current = _lease_now(now)
root = _ensure_lock_dir(lock_dir)
path = lock_file_path(
remote=remote, org=org, repo=repo, issue_number=issue_number, lock_dir=root
)
sentinel = flock_path(path)
try:
with _exclusive_file_lock(sentinel):
lock = read_lock_file(path)
if not lock:
return _refusal([f"no durable lock for issue #{issue_number}"])
if not is_legacy_lease(lock):
return _refusal(
[
"lock is already on the heartbeat lifecycle; use the "
"heartbeat path"
],
lifecycle=lease_lifecycle_version(lock),
legacy_lease=False,
)
reasons = _ownership_refusals(
lock,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=worktree_path,
identity=identity,
profile=profile,
)
if reasons:
return _refusal(reasons)
current_generation = lock_generation(lock)
if (
expected_generation is not None
and current_generation != expected_generation
):
return _refusal(
[
f"lock generation changed: expected {expected_generation}, "
f"found {current_generation} (fail closed)"
],
lock_generation=current_generation,
)
freshness = assess_lock_freshness(lock, now=current)
if not freshness.get("live"):
return _refusal(
[
f"legacy lease is not live ({freshness.get('status')}): "
f"{freshness.get('reason')}; rebinding is not a recovery "
"path for a lapsed lease"
],
freshness=freshness,
)
policy = lease_policy.policy_for(lease_task_class(lock))
expires = current + timedelta(minutes=policy.initial_ttl_minutes)
session_id = mint_task_session_id(lease_task_class(lock))
record = dict(lock)
lease = dict(record.get("work_lease") or {})
legacy_origin = {
"created_at": lease.get("created_at"),
"expires_at": lease.get("expires_at"),
"last_heartbeat_at": lease.get("last_heartbeat_at"),
"lifecycle": lease_policy.LIFECYCLE_LEGACY,
}
lease["lifecycle_version"] = lease_policy.LIFECYCLE_HEARTBEAT_V1
lease["task_session_id"] = session_id
lease["created_at"] = _format_lease_timestamp(current)
lease["last_heartbeat_at"] = _format_lease_timestamp(current)
lease["expires_at"] = _format_lease_timestamp(expires)
lease["heartbeat_count"] = 1
record["work_lease"] = lease
record["legacy_rebind"] = {
"rebound_at": _format_lease_timestamp(current),
"task_session_id": session_id,
"prior_generation": current_generation,
"legacy_origin": legacy_origin,
"reason": (
"legacy lock rebound into the heartbeat lifecycle by its exact "
"recorded owner"
),
}
record["lock_generation"] = current_generation + 1
save_lock_file(path, record)
except LockContentionError as exc:
return _refusal([f"issue #{issue_number} lock contention: {exc} (fail closed)"])
return {
"success": True,
"performed": True,
"issue_number": issue_number,
"task_session_id": session_id,
"lock_generation": record["lock_generation"],
"prior_generation": current_generation,
"lifecycle": lease_policy.LIFECYCLE_HEARTBEAT_V1,
"legacy_rebind": record["legacy_rebind"],
"expires_at": lease["expires_at"],
"last_heartbeat_at": lease["last_heartbeat_at"],
"lock_file_path": path,
"freshness": assess_lock_freshness(record, now=current),
}
def read_session_issue_lock(lock_dir: str | None = None) -> dict[str, Any] | None:
root = (lock_dir or default_lock_dir()).strip()
pointer = read_lock_file(session_pointer_path(root))
@@ -336,6 +673,16 @@ def _parse_lease_timestamp(value: str | None) -> datetime | None:
return None
def _format_lease_timestamp(value: datetime) -> str:
"""Serialize a lease timestamp in the durable ``...Z`` form already on disk."""
return (
value.astimezone(timezone.utc)
.replace(microsecond=0)
.isoformat()
.replace("+00:00", "Z")
)
def lease_expires_at(lock: dict[str, Any] | None) -> datetime | None:
if not lock:
return None
@@ -356,60 +703,216 @@ def is_lease_live(lock: dict[str, Any] | None, *, now: datetime | None = None) -
return assess_lock_freshness(lock, now=now)["live"]
def lease_task_class(lock_data: dict[str, Any] | None) -> str:
"""Policy task class for a durable lock; author work when unrecorded."""
lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None
if isinstance(lease, dict):
recorded = str(lease.get("operation_type") or "").strip()
if recorded:
return recorded
return AUTHOR_ISSUE_WORK_LEASE
def lease_lifecycle_version(lock_data: dict[str, Any] | None) -> str:
"""Read the durable lifecycle marker (#790 AC-N8).
The marker is the *only* discriminator between a heartbeat-lifecycle lease
and a legacy one. Timestamps are deliberately not consulted: a lock minted
before this lifecycle existed has ``last_heartbeat_at == created_at``
forever, and reading that equality as "recently heartbeated" would treat
every never-heartbeated legacy lock as fresh — the precise inversion AC-N8
forbids. A newly minted heartbeat lease also has the two equal, so the
equality carries no information in either direction.
"""
lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None
if isinstance(lease, dict):
recorded = str(lease.get("lifecycle_version") or "").strip()
if recorded:
return recorded
return lease_policy.LIFECYCLE_LEGACY
def is_legacy_lease(lock_data: dict[str, Any] | None) -> bool:
"""True when a lock predates the shared heartbeat lifecycle."""
return lease_lifecycle_version(lock_data) != lease_policy.LIFECYCLE_HEARTBEAT_V1
def lease_task_session_id(lock_data: dict[str, Any] | None) -> str:
"""Recorded per-task session identifier, or empty for a legacy lock."""
lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None
if isinstance(lease, dict):
return str(lease.get("task_session_id") or "").strip()
return ""
def mint_task_session_id(task_class: str = AUTHOR_ISSUE_WORK_LEASE) -> str:
"""Mint an ownership key for one task (#790 AC-N1).
Deliberately contains no process identifier. The recorded PID belongs to the
long-lived MCP daemon, which outlives any individual task and is reused by
every task it serves, so PID digits cannot identify *which* task holds a
claim. The PID is still recorded alongside this value as evidence.
"""
prefix = _sanitize_segment(str(task_class or AUTHOR_ISSUE_WORK_LEASE))
return f"{prefix}-{uuid.uuid4().hex[:16]}"
def _lease_heartbeat_at(lock_data: dict[str, Any] | None) -> datetime | None:
lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None
heartbeat_at = None
if isinstance(lock_data, dict):
heartbeat_at = _parse_lease_timestamp(lock_data.get("last_heartbeat_at"))
if heartbeat_at is None and isinstance(lease, dict):
heartbeat_at = _parse_lease_timestamp(lease.get("last_heartbeat_at"))
return heartbeat_at
def assess_lock_freshness(
lock_data: dict[str, Any] | None,
*,
now: datetime | None = None,
) -> dict[str, Any]:
"""Classify a lock as live, expired, stale, or absent."""
"""Classify a lock as live, expired, stale, or absent.
#790 Slice A makes the heartbeat load-bearing. Before this change
``last_heartbeat_at`` was parsed and then never consulted: liveness was
decided entirely by the absolute ``expires_at`` and by PID liveness, and
since the recorded PID is the long-lived MCP daemon, an abandoned author
task stayed "live" for the full four-hour TTL.
Two rules govern the rewrite:
* **An alive PID never establishes freshness** (AC-N2). It proves the daemon
is up, nothing about the task. It is recorded as evidence and no branch
returns ``live`` because of it.
* **A dead PID still corroborates staleness.** The dead-PID band is
unchanged and still precedes every heartbeat evaluation, so #753
dead-session recovery keys on exactly the classification it always did.
Legacy leases (AC-N8) keep their recorded absolute expiry and are never
evaluated against the short heartbeat grace, so deploying this change cannot
make an existing claim instantly reclaimable.
"""
current = _lease_now(now)
if not lock_data:
return {
"status": "absent",
"status": STATUS_ABSENT,
"live": False,
"stale": False,
"reason": "no lock record",
}
expires_at = lease_expires_at(lock_data)
lease = lock_data.get("work_lease")
heartbeat_at = _parse_lease_timestamp(lock_data.get("last_heartbeat_at"))
if heartbeat_at is None and isinstance(lease, dict):
heartbeat_at = _parse_lease_timestamp(lease.get("last_heartbeat_at"))
expires_at = lease_expires_at(lock_data)
heartbeat_at = _lease_heartbeat_at(lock_data)
created_at = (
_parse_lease_timestamp(lease.get("created_at"))
if isinstance(lease, dict)
else None
)
pid = lock_data.get("session_pid")
if pid is None:
pid = lock_data.get("pid")
# Evidence only. Never consulted to grant liveness (AC-N2).
pid_alive = is_process_alive(pid) if pid is not None else False
if expires_at and expires_at <= current:
return {
"status": "expired",
"live": False,
"stale": True,
"reason": f"lease expired at {expires_at.isoformat()}",
"pid_alive": pid_alive,
}
lifecycle = lease_lifecycle_version(lock_data)
legacy = lifecycle != lease_policy.LIFECYCLE_HEARTBEAT_V1
policy = lease_policy.policy_for(lease_task_class(lock_data))
if pid is not None and not pid_alive:
return {
"status": "stale",
"live": False,
"stale": True,
"reason": f"owner pid {pid} is not alive",
"pid_alive": False,
}
return {
"status": "live",
"live": True,
"stale": False,
"reason": "lock heartbeat and lease are fresh",
evidence: dict[str, Any] = {
"pid_alive": pid_alive,
"lifecycle": lifecycle,
"legacy_lease": legacy,
"task_session_id": lease_task_session_id(lock_data) or None,
"heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None,
"expires_at": expires_at.isoformat() if expires_at else None,
}
def _result(status: str, *, live: bool, reason: str, **extra: Any) -> dict[str, Any]:
return {
"status": status,
"live": live,
"stale": not live and status != STATUS_ABSENT,
"reason": reason,
**evidence,
**extra,
}
if legacy:
# AC-N8: the preserved absolute expiry is the only clock for a lock
# written before task-session heartbeats existed.
if expires_at and expires_at <= current:
return _result(
STATUS_EXPIRED,
live=False,
reason=f"lease expired at {expires_at.isoformat()}",
)
if pid is not None and not pid_alive:
return _result(
STATUS_STALE, live=False, reason=f"owner pid {pid} is not alive"
)
return _result(
STATUS_LIVE,
live=True,
reason=(
"legacy lease is within its recorded absolute expiry; the "
"heartbeat grace does not apply retroactively"
),
legacy_expiry_preserved=True,
)
# ── Heartbeat lifecycle ──
if pid is not None and not pid_alive:
# Unchanged dead-PID band: #753 recovery depends on this exact status.
return _result(STATUS_STALE, live=False, reason=f"owner pid {pid} is not alive")
if heartbeat_at is None:
# Contradictory: a heartbeat lease must carry a heartbeat. Fail closed.
return _result(
STATUS_STALE_MISSED_HEARTBEAT,
live=False,
reason=(
f"lease declares lifecycle '{lifecycle}' but records no "
"last_heartbeat_at (fail closed)"
),
)
if policy.absolute_cap_hours and created_at is not None:
cap_at = created_at + timedelta(hours=policy.absolute_cap_hours)
if cap_at <= current:
return _result(
STATUS_STALE_ABSOLUTE_CAP,
live=False,
reason=(
f"lease exceeded its {policy.absolute_cap_hours}h absolute cap "
f"at {cap_at.isoformat()}; canonical re-adoption is required"
),
absolute_cap_at=cap_at.isoformat(),
)
grace_at = heartbeat_at + timedelta(minutes=policy.missed_heartbeat_grace_minutes)
if grace_at <= current or (expires_at is not None and expires_at <= current):
return _result(
STATUS_STALE_MISSED_HEARTBEAT,
live=False,
reason=(
f"no valid heartbeat since {heartbeat_at.isoformat()}; the "
f"{policy.missed_heartbeat_grace_minutes}min grace lapsed at "
f"{grace_at.isoformat()}"
),
missed_heartbeat_since=grace_at.isoformat(),
)
warning_at = heartbeat_at + timedelta(minutes=policy.stale_warning_minutes)
return _result(
STATUS_LIVE,
live=True,
reason="lease heartbeat is fresh within the configured grace",
heartbeat_warning=warning_at <= current,
)
def _same_realpath(left: str | None, right: str | None) -> bool:
if not left or not right:
@@ -446,6 +949,27 @@ def assess_expired_lock_reclaim(
"reasons": ["lock is still live; cannot reclaim (fail closed)"],
"freshness": freshness,
}
status = str(freshness.get("status") or "")
if status in (STATUS_STALE_MISSED_HEARTBEAT, STATUS_STALE_ABSOLUTE_CAP):
# #790: under the heartbeat lifecycle the heartbeat *is* the liveness
# proof, so a session that stopped heartbeating past its grace has
# released its claim by definition. Requiring a dead PID on top of that
# would reinstate the original defect — the recorded PID is the shared
# daemon, which stays alive across every abandoned task it ever served.
#
# This band is unreachable for a legacy lease (AC-N8), so no lock
# written before this lifecycle can be reclaimed by this path.
return {
"reclaim_allowed": True,
"reasons": [
f"heartbeat-lifecycle lease is {status}: {freshness.get('reason')}"
],
"freshness": freshness,
"prior_branch": existing_lock.get("branch_name"),
"prior_worktree": existing_lock.get("worktree_path"),
"prior_pid": existing_lock.get("session_pid") or existing_lock.get("pid"),
"prior_task_session_id": lease_task_session_id(existing_lock) or None,
}
pid = existing_lock.get("session_pid")
if pid is None:
pid = existing_lock.get("pid")
+212
View File
@@ -0,0 +1,212 @@
"""Central lease policy configuration (#790 Slice A, AC-N7).
The single authoritative source for every lease duration in the project. Before
this module the numbers were scattered: a four-hour author TTL was declared
twice (``issue_lock_store`` and ``gitea_mcp_server``), the reviewer/merger
sliding window lived in ``reviewer_pr_lease``, the conflict-fix window in
``pr_work_lease``, and the control-plane default in ``control_plane_db``.
Nothing tied them together, so tuning one class silently diverged from the
others and no reader could answer "how long does a lease live?" without
grepping four files.
AC-N7 requires that this configuration exist *before* the first heartbeat and
TTL behavior that reads from it, so it ships in Slice A rather than trailing the
code it governs.
Deliberate boundaries:
* **Declaration is not rewiring.** Every task class is declared here, but only
those with ``heartbeat_lifecycle_active`` were migrated onto the shared
heartbeat lifecycle in Slice A — currently ``author_issue_work`` alone.
Reviewer, merger, and conflict-fix leases keep their own existing behavior
until Slice C moves them; their numbers are recorded here so the two cannot
drift apart unnoticed, and ``tests/test_issue_790_lease_policy.py`` asserts
the recorded values still equal the constants those modules use.
* **No policy decision lives here.** This module answers "how long", never "may
this session proceed". Freshness, reclaim, and renewal dispositions stay in
``issue_lock_store``.
"""
from __future__ import annotations
import os
from dataclasses import dataclass
from typing import Any
# Task classes. Only the first is migrated onto the shared lifecycle in Slice A.
TASK_CLASS_AUTHOR_ISSUE_WORK = "author_issue_work"
TASK_CLASS_REVIEWER_PR = "reviewer_pr"
TASK_CLASS_MERGER_PR = "merger_pr"
TASK_CLASS_CONFLICT_FIX = "conflict_fix"
# Durable marker for a lease minted under the shared heartbeat lifecycle.
#
# #790 AC-N8: this explicit marker — never a timestamp comparison — is what
# distinguishes a heartbeat-lifecycle lease from a legacy one. A lock written
# before this lifecycle existed carries no marker and reads as
# ``LIFECYCLE_LEGACY``.
LIFECYCLE_HEARTBEAT_V1 = "heartbeat-v1"
LIFECYCLE_LEGACY = "legacy"
_ENV_PREFIX = "GITEA_LEASE_POLICY"
@dataclass(frozen=True)
class LeasePolicy:
"""Durations governing one task class.
All intervals are minutes except ``absolute_cap_hours``. ``None`` for the
cap means the class has no maximum continuous duration.
"""
task_class: str
initial_ttl_minutes: float
heartbeat_cadence_minutes: float
stale_warning_minutes: float
missed_heartbeat_grace_minutes: float
absolute_cap_hours: float | None
recovery_grace_minutes: float
terminal_race_drain_minutes: float
terminal_retirement_eligible: bool
heartbeat_lifecycle_active: bool
# Defaults. ``author_issue_work`` adopts the reviewer window proven by #747
# rather than inventing new numbers: a lease expires 10 minutes after its last
# valid heartbeat, warns at half that, and an actively heartbeating session is
# never evicted. The prior value was a fixed four hours (240 minutes) that no
# heartbeat could shorten — the defect this issue exists to correct.
_DEFAULTS: dict[str, LeasePolicy] = {
TASK_CLASS_AUTHOR_ISSUE_WORK: LeasePolicy(
task_class=TASK_CLASS_AUTHOR_ISSUE_WORK,
initial_ttl_minutes=10.0,
heartbeat_cadence_minutes=2.0,
stale_warning_minutes=5.0,
missed_heartbeat_grace_minutes=10.0,
absolute_cap_hours=8.0,
recovery_grace_minutes=10.0,
terminal_race_drain_minutes=2.0,
terminal_retirement_eligible=True,
heartbeat_lifecycle_active=True,
),
# Declared, not rewired. These mirror reviewer_pr_lease.LEASE_TTL_MINUTES
# and STALE_WARNING_MINUTES; Slice C migrates the call sites.
TASK_CLASS_REVIEWER_PR: LeasePolicy(
task_class=TASK_CLASS_REVIEWER_PR,
initial_ttl_minutes=10.0,
heartbeat_cadence_minutes=2.0,
stale_warning_minutes=5.0,
missed_heartbeat_grace_minutes=10.0,
absolute_cap_hours=None,
recovery_grace_minutes=10.0,
terminal_race_drain_minutes=2.0,
terminal_retirement_eligible=False,
heartbeat_lifecycle_active=False,
),
TASK_CLASS_MERGER_PR: LeasePolicy(
task_class=TASK_CLASS_MERGER_PR,
initial_ttl_minutes=10.0,
heartbeat_cadence_minutes=2.0,
stale_warning_minutes=5.0,
missed_heartbeat_grace_minutes=10.0,
absolute_cap_hours=None,
recovery_grace_minutes=10.0,
terminal_race_drain_minutes=2.0,
terminal_retirement_eligible=False,
heartbeat_lifecycle_active=False,
),
# Mirrors pr_work_lease.DEFAULT_CONFLICT_FIX_TTL_MINUTES. Deliberately left
# at its current window; shortening it is Slice C's call, not this slice's.
TASK_CLASS_CONFLICT_FIX: LeasePolicy(
task_class=TASK_CLASS_CONFLICT_FIX,
initial_ttl_minutes=120.0,
heartbeat_cadence_minutes=2.0,
stale_warning_minutes=5.0,
missed_heartbeat_grace_minutes=10.0,
absolute_cap_hours=None,
recovery_grace_minutes=10.0,
terminal_race_drain_minutes=2.0,
terminal_retirement_eligible=False,
heartbeat_lifecycle_active=False,
),
}
_NUMERIC_FIELDS = (
"initial_ttl_minutes",
"heartbeat_cadence_minutes",
"stale_warning_minutes",
"missed_heartbeat_grace_minutes",
"absolute_cap_hours",
"recovery_grace_minutes",
"terminal_race_drain_minutes",
)
def env_var_name(task_class: str, field: str) -> str:
"""Environment variable that overrides one field of one task class."""
return f"{_ENV_PREFIX}_{task_class.upper()}_{field.upper()}"
def _override(task_class: str, field: str, default: float | None) -> float | None:
"""Read one override, falling back to *default* on anything unusable.
A malformed or non-positive override is ignored rather than raised: a typo
in an environment variable must not be able to mint a zero-length lease that
makes every claim instantly reclaimable, nor crash the server at import.
"""
raw = (os.environ.get(env_var_name(task_class, field)) or "").strip()
if not raw:
return default
try:
value = float(raw)
except (TypeError, ValueError):
return default
if value <= 0:
return default
return value
def policy_for(task_class: str) -> LeasePolicy:
"""Return the effective policy for *task_class*.
Unknown task classes fall back to the author policy, which is the most
conservative migrated class, rather than raising — a new caller must never
be able to crash a lock write by naming a class this table has not learned.
"""
key = str(task_class or "").strip() or TASK_CLASS_AUTHOR_ISSUE_WORK
base = _DEFAULTS.get(key) or _DEFAULTS[TASK_CLASS_AUTHOR_ISSUE_WORK]
resolved = {
field: _override(base.task_class, field, getattr(base, field))
for field in _NUMERIC_FIELDS
}
if all(resolved[field] == getattr(base, field) for field in _NUMERIC_FIELDS):
return base
return LeasePolicy(
task_class=base.task_class,
terminal_retirement_eligible=base.terminal_retirement_eligible,
heartbeat_lifecycle_active=base.heartbeat_lifecycle_active,
**resolved,
)
def known_task_classes() -> tuple[str, ...]:
"""Every declared task class, migrated or not."""
return tuple(_DEFAULTS)
def describe(task_class: str) -> dict[str, Any]:
"""Serializable view of a policy, for audit records and tool payloads."""
policy = policy_for(task_class)
return {
"task_class": policy.task_class,
"initial_ttl_minutes": policy.initial_ttl_minutes,
"heartbeat_cadence_minutes": policy.heartbeat_cadence_minutes,
"stale_warning_minutes": policy.stale_warning_minutes,
"missed_heartbeat_grace_minutes": policy.missed_heartbeat_grace_minutes,
"absolute_cap_hours": policy.absolute_cap_hours,
"recovery_grace_minutes": policy.recovery_grace_minutes,
"terminal_race_drain_minutes": policy.terminal_race_drain_minutes,
"terminal_retirement_eligible": policy.terminal_retirement_eligible,
"heartbeat_lifecycle_active": policy.heartbeat_lifecycle_active,
"lifecycle_version": LIFECYCLE_HEARTBEAT_V1,
}
+6 -8
View File
@@ -32,14 +32,12 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"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).
"rebind_dirty_same_claimant_author_session": {
"permission": "gitea.issue.comment",
"role": "author",
},
"gitea_rebind_dirty_same_claimant_author_session": {
# #790 Slice A: prove an owned author lease is still active. Strictly
# narrower than lock_issue — it can only slide a lease this exact session
# already owns, never acquire, take over, or revive one — so it gates on the
# same authority rather than introducing an operation name that every
# already-configured author profile would be missing.
"heartbeat_issue_lock": {
"permission": "gitea.issue.comment",
"role": "author",
},
-352
View File
@@ -1639,358 +1639,6 @@ 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()
@@ -1,845 +0,0 @@
"""Integration tests for dirty same-claimant author-session rebind (#864).
Uses real temp git repos/worktrees and a temp GITEA_ISSUE_LOCK_DIR. Does not
mutate any real #860/#864 worktree on disk.
"""
from __future__ import annotations
import json
import os
import subprocess
import sys
import tempfile
from datetime import datetime, timedelta, timezone
from pathlib import Path
import pytest
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
import dirty_same_claimant_session_rebind as rebind # noqa: E402
import issue_lock_provenance # noqa: E402
import issue_lock_store as ils # noqa: E402
import issue_lock_worktree # noqa: E402
ISSUE = 864
BRANCH = f"fix/issue-{ISSUE}-dirty-same-claimant-session-rebind"
REMOTE = "prgs"
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
IDENTITY = "jcwalker3"
PROFILE = "prgs-author"
def _git(cwd: str, *args: str, check: bool = True) -> subprocess.CompletedProcess:
return subprocess.run(
["git", "-C", cwd, *args],
capture_output=True,
text=True,
check=check,
)
def dead_pid() -> int:
proc = subprocess.Popen([sys.executable, "-c", "pass"])
proc.wait()
return proc.pid
def future_ts(hours: int = 4) -> str:
return (
(datetime.now(timezone.utc) + timedelta(hours=hours))
.isoformat()
.replace("+00:00", "Z")
)
@pytest.fixture
def lock_dir(tmp_path, monkeypatch):
d = tmp_path / "issue-locks"
d.mkdir()
monkeypatch.setenv("GITEA_ISSUE_LOCK_DIR", str(d))
return str(d)
@pytest.fixture
def dirty_repo(tmp_path):
"""Canonical repo root with branches/<name> worktree and dirty content."""
root = tmp_path / "repo"
root.mkdir()
main = root / "main"
main.mkdir()
subprocess.run(["git", "init", "-q", str(main)], check=True, capture_output=True)
_git(str(main), "config", "user.email", "t@t")
_git(str(main), "config", "user.name", "t")
(main / "README.md").write_text("base\n", encoding="utf-8")
_git(str(main), "add", "README.md")
_git(str(main), "commit", "-q", "-m", "base")
_git(str(main), "branch", "-M", "master")
# Bare remote + origin tracking so remote head is observable offline.
bare = tmp_path / "remote.git"
subprocess.run(
["git", "init", "--bare", "-q", str(bare)], check=True, capture_output=True
)
_git(str(main), "remote", "add", "origin", str(bare))
_git(str(main), "push", "-q", "origin", "master:master")
branches = root / "branches"
branches.mkdir()
wt_name = f"fix-issue-{ISSUE}-dirty-same-claimant-session-rebind"
wt = branches / wt_name
_git(str(main), "worktree", "add", "-q", "-b", BRANCH, str(wt))
_git(str(wt), "push", "-q", "-u", "origin", BRANCH)
# Seed committed files we will dirty.
tracked = [
"dirty_same_claimant_session_rebind.py",
"issue_lock_provenance.py",
"task_capability_map.py",
]
for rel in tracked:
p = wt / rel
p.parent.mkdir(parents=True, exist_ok=True)
p.write_text(f"seed {rel}\n", encoding="utf-8")
_git(str(wt), "add", *tracked)
_git(str(wt), "commit", "-q", "-m", "seed tracked")
_git(str(wt), "push", "-q", "origin", BRANCH)
# Dirty tracked + untracked.
for rel in tracked:
(wt / rel).write_text(f"dirty {rel}\n", encoding="utf-8")
untracked = [
"tests/test_dirty_same_claimant_session_rebind.py",
"docs/runbook-dirty-rebind.md",
"scratch/notes-untracked.txt",
"extra_untracked.txt",
]
for rel in untracked:
p = wt / rel
p.parent.mkdir(parents=True, exist_ok=True)
p.write_text(f"untracked {rel}\n", encoding="utf-8")
inv = rebind.collect_dirty_inventory(str(wt))
assert inv["ok"], inv.get("reasons")
head = _git(str(wt), "rev-parse", "HEAD").stdout.strip()
remote_head = _git(
str(wt), "rev-parse", f"refs/remotes/origin/{BRANCH}"
).stdout.strip()
assert head == remote_head
return {
"root": str(root),
"main": str(main),
"worktree": str(wt),
"branch": BRANCH,
"inventory": inv,
"local_head": head,
"remote_head": remote_head,
"dirty_paths": list(inv["dirty_paths"]),
"fingerprints": dict(inv["fingerprints"]),
}
def _make_lock(
*,
worktree: str,
pid: int,
lock_dir: str,
identity: str = IDENTITY,
profile: str = PROFILE,
**overrides,
) -> dict:
lease = {
"operation_type": ils.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"branch": BRANCH,
"worktree_path": worktree,
"claimant": {"username": identity, "profile": profile},
"created_at": "2026-01-01T00:00:00Z",
"expires_at": future_ts(),
"last_heartbeat_at": "2026-01-01T00:00:00Z",
}
lock = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": worktree,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"session_pid": pid,
"pid": pid,
"work_lease": lease,
"lock_generation": 1,
"lock_provenance": issue_lock_provenance.build_sanctioned_lock_provenance(
tool="gitea_lock_issue",
claimant={"username": identity, "profile": profile},
),
}
lock.update(overrides)
path = ils.lock_file_path(
remote=REMOTE, org=ORG, repo=REPO, issue_number=ISSUE, lock_dir=lock_dir
)
lock["lock_file_path"] = path
ils.save_lock_file(path, lock)
# Stale session pointer for the dead owner.
ptr = {
"pid": pid,
"lock_file_path": path,
"issue_number": ISSUE,
"branch_name": BRANCH,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
}
ils.save_lock_file(os.path.join(lock_dir, f"session-{pid}.json"), ptr)
return ils.read_lock_file(path) or lock
def _apply_kwargs(repo, lock, lock_dir, **overrides):
kwargs = {
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": repo["worktree"],
"claimant_identity": IDENTITY,
"claimant_profile": PROFILE,
"old_pid": lock.get("session_pid") or lock.get("pid"),
"expected_local_head": repo["local_head"],
"expected_remote_head": repo["remote_head"],
"expected_dirty_paths": repo["dirty_paths"],
"expected_fingerprints": repo["fingerprints"],
"existing_lock": lock,
"current_identity": IDENTITY,
"current_profile": PROFILE,
"role_kind": "author",
"current_pid": os.getpid(),
"current_branch": BRANCH,
"local_head": repo["local_head"],
"remote_head": repo["remote_head"],
"dirty_inventory": repo["inventory"],
"competing_live_locks": [],
"competing_sessions": [],
"workflow_lease_active": False,
"repo_root": repo["root"],
"dry_run": False,
"lock_dir": lock_dir,
}
kwargs.update(overrides)
return kwargs
# ── 1. Successful dead-PID same-claimant dirty rebind ───────────────────────
def test_successful_dead_pid_same_claimant_dirty_rebind(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert result["success"], result
assert result["outcome"] == rebind.REBIND_SANCTIONED
assert result["old_pid"] == old
assert result["new_pid"] == os.getpid()
assert result["generation_after"] == result["generation_before"] + 1
rebound = ils.read_lock_file(result["lock_path"])
assert rebound is not None
assert int(rebound["session_pid"]) == os.getpid()
assert int(rebound["pid"]) == os.getpid()
assert (
rebound.get("lock_provenance", {}).get("source")
== issue_lock_provenance.SOURCE_DIRTY_SAME_CLAIMANT_REBIND
)
assert rebound.get("rebind_record", {}).get("old_pid") == old
# ── 2. Byte-for-byte preservation ───────────────────────────────────────────
def test_byte_for_byte_preservation(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
before = {
rel: rebind.content_fingerprint(os.path.join(dirty_repo["worktree"], rel))
for rel in dirty_repo["dirty_paths"]
}
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert result["success"], result
after = {
rel: rebind.content_fingerprint(os.path.join(dirty_repo["worktree"], rel))
for rel in dirty_repo["dirty_paths"]
}
assert before == after
assert result["fingerprints"] == before
# ── 3. Exact dirty-path and fingerprint enforcement ─────────────────────────
def test_extra_dirty_path_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
pins = list(dirty_repo["dirty_paths"])[:-1] # missing one observed path
fps = {p: dirty_repo["fingerprints"][p] for p in pins}
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
expected_dirty_paths=pins,
expected_fingerprints=fps,
)
)
assert not result["success"]
assert any("unexpected paths" in r for r in result["reasons"])
def test_missing_expected_dirty_path_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
pins = list(dirty_repo["dirty_paths"]) + ["not_really_dirty.txt"]
fps = dict(dirty_repo["fingerprints"])
fps["not_really_dirty.txt"] = "0" * 64
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
expected_dirty_paths=pins,
expected_fingerprints=fps,
)
)
assert not result["success"]
assert any("missing expected" in r for r in result["reasons"])
def test_modified_fingerprint_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
fps = dict(dirty_repo["fingerprints"])
victim = dirty_repo["dirty_paths"][0]
fps[victim] = "f" * 64
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
expected_fingerprints=fps,
)
)
assert not result["success"]
assert any("fingerprint disagreement" in r for r in result["reasons"])
# ── 4. Atomic session-pointer replacement ───────────────────────────────────
def test_session_pointer_points_to_lock(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert result["success"], result
new_ptr_path = os.path.join(lock_dir, f"session-{os.getpid()}.json")
assert os.path.exists(new_ptr_path)
ptr = ils.read_lock_file(new_ptr_path)
assert ptr is not None
assert os.path.realpath(ptr["lock_file_path"]) == os.path.realpath(
result["lock_path"]
)
# Old pointer removed when it targeted this lock.
old_ptr = os.path.join(lock_dir, f"session-{old}.json")
assert not os.path.exists(old_ptr)
assert result.get("removed_old_session_pointer") is True
# ── 5. Retry after interruption (journal mid-state) ─────────────────────────
def test_retry_after_journal_mid_state(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
jpath = rebind.journal_path(lock_dir, ISSUE)
rebind._atomic_write_json(
jpath,
{
"phase": rebind.JOURNAL_PHASE_PRE_BIND,
"issue_number": ISSUE,
"old_pid": old,
"new_pid": os.getpid(),
"expected_generation": 1,
},
)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert result["success"], result
assert result["journal_phase"] == rebind.JOURNAL_PHASE_COMPLETE
# Second apply is already_rebound (retry-safe).
rebound_lock = ils.read_lock_file(result["lock_path"])
result2 = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
rebound_lock,
lock_dir,
old_pid=old,
existing_lock=rebound_lock,
)
)
assert result2["success"], result2
assert result2["already_rebound"] is True
# ── 6. Active-PID refusal ───────────────────────────────────────────────────
def test_active_pid_refused(dirty_repo, lock_dir):
live = os.getpid()
# Use a different "current" identity of session via fake current_pid...
# Owner is live (this process). Rebind must refuse.
lock = _make_lock(worktree=dirty_repo["worktree"], pid=live, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=live,
current_pid=live + 10_000_000, # distinct "new" session id for pin check
)
)
assert not result["success"]
assert any("still alive" in r for r in result["reasons"])
# ── 7. Foreign claimant refusal ─────────────────────────────────────────────
def test_foreign_claimant_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
current_identity="someone-else",
claimant_identity="someone-else",
)
)
assert not result["success"]
assert any("foreign claimant" in r or "does not match" in r for r in result["reasons"])
# ── 8. Profile mismatch refusal ─────────────────────────────────────────────
def test_profile_mismatch_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
current_profile="other-profile",
claimant_profile="other-profile",
)
)
assert not result["success"]
assert any("profile" in r for r in result["reasons"])
# ── 9. Competing session/lock/lease refusal ─────────────────────────────────
def test_competing_live_lock_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
competing = [
{
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": dirty_repo["worktree"] + "-other",
"pid": os.getpid(),
}
]
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
competing_live_locks=competing,
)
)
assert not result["success"]
assert any("competing live lock" in r for r in result["reasons"])
def test_competing_session_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
competing_sessions=[
{"pid": os.getpid(), "lock_file_path": lock["lock_file_path"], "live": True}
],
)
)
# current_pid is os.getpid(), so same session is skipped — use another live pid.
# Spawn a long-lived process to act as competing live session.
rival = subprocess.Popen([sys.executable, "-c", "import time; time.sleep(30)"])
try:
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
competing_sessions=[
{
"pid": rival.pid,
"lock_file_path": lock["lock_file_path"],
"live": True,
}
],
)
)
assert not result["success"]
assert any("competing live session" in r for r in result["reasons"])
finally:
rival.kill()
rival.wait()
def test_workflow_lease_active_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
workflow_lease_active=True,
)
)
assert not result["success"]
assert any("workflow lease" in r for r in result["reasons"])
# ── 10. Local- and remote-head movement refusal ─────────────────────────────
def test_local_head_movement_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
expected_local_head="b" * 40,
)
)
assert not result["success"]
assert any("local head" in r for r in result["reasons"])
def test_remote_head_movement_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
expected_remote_head="c" * 40,
)
)
assert not result["success"]
assert any("remote head" in r for r in result["reasons"])
# ── 11. Path/symlink/registration mismatches ────────────────────────────────
def test_worktree_not_under_branches_refused(dirty_repo, lock_dir, tmp_path):
old = dead_pid()
# Use a path outside branches/ as the declared worktree (still real dir).
outside = tmp_path / "outside-wt"
outside.mkdir()
lock = _make_lock(worktree=str(outside), pid=old, lock_dir=lock_dir)
# Inventory empty for outside path; use empty pins to hit path gate first
# by providing matching empty-ish inventory after we force path checks.
inv = {
"dirty_paths": dirty_repo["dirty_paths"],
"fingerprints": dirty_repo["fingerprints"],
"ok": True,
"reasons": [],
}
result = rebind.assess_dirty_same_claimant_session_rebind(
remote=REMOTE,
org=ORG,
repo=REPO,
issue_number=ISSUE,
branch_name=BRANCH,
worktree_path=str(outside),
claimant_identity=IDENTITY,
claimant_profile=PROFILE,
old_pid=old,
expected_local_head=dirty_repo["local_head"],
expected_remote_head=dirty_repo["remote_head"],
expected_dirty_paths=dirty_repo["dirty_paths"],
expected_fingerprints=dirty_repo["fingerprints"],
existing_lock=lock,
current_identity=IDENTITY,
current_profile=PROFILE,
role_kind="author",
current_pid=os.getpid(),
current_branch=BRANCH,
local_head=dirty_repo["local_head"],
remote_head=dirty_repo["remote_head"],
dirty_inventory=inv,
repo_root=dirty_repo["root"],
)
assert not result["rebind_sanctioned"]
assert any("branches/" in r for r in result["reasons"])
def test_lock_worktree_mismatch_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(
worktree=dirty_repo["worktree"] + "-elsewhere",
pid=old,
lock_dir=lock_dir,
)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert not result["success"]
assert any("does not match declared" in r for r in result["reasons"])
# ── 12. Malformed lock/session records ──────────────────────────────────────
def test_malformed_lock_missing_pid_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
lock.pop("session_pid", None)
lock.pop("pid", None)
ils.save_lock_file(lock["lock_file_path"], lock)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert not result["success"]
assert any("incomplete" in r or "session_pid" in r for r in result["reasons"])
def test_empty_old_pid_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=None)
)
assert not result["success"]
assert any("old_pid" in r for r in result["reasons"])
def test_reviewer_role_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old, role_kind="reviewer")
)
assert not result["success"]
assert any("reviewer" in r for r in result["reasons"])
def test_reconciler_without_authorize_refused(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(
dirty_repo,
lock,
lock_dir,
old_pid=old,
role_kind="reconciler",
authorize_reconciler_execute=False,
)
)
assert not result["success"]
assert any("authorize_reconciler_execute" in r for r in result["reasons"])
# ── 13. No duplicate ownership after success or retry ───────────────────────
def test_no_duplicate_ownership_after_success_or_retry(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
r1 = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old)
)
assert r1["success"], r1
rebound = ils.read_lock_file(r1["lock_path"])
r2 = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, rebound, lock_dir, old_pid=old, existing_lock=rebound)
)
assert r2["success"], r2
assert r2["already_rebound"] is True
# Only one durable lock file for this issue; session pointer is current pid.
# Skip session pointers and rebind journals (dotfiles / non-lock records).
matching = []
for p in ils.iter_lock_files(lock_dir):
name = os.path.basename(p)
if name.startswith(".") or name.startswith("session-"):
continue
rec = ils.read_lock_file(p)
if not rec:
continue
if (
rec.get("issue_number") == ISSUE
and rec.get("remote") == REMOTE
and rec.get("branch_name") == BRANCH
and rec.get("session_pid") is not None
):
matching.append(rec)
assert len(matching) == 1
assert int(matching[0]["session_pid"]) == os.getpid()
# No live session pointer for the dead old pid.
assert not os.path.exists(os.path.join(lock_dir, f"session-{old}.json"))
# ── 14. Ordinary dirty-worktree locking remains fail-closed ─────────────────
def test_ordinary_dirty_lock_worktree_assessment_blocks(dirty_repo):
porcelain = dirty_repo["inventory"]["porcelain_status"]
assessment = issue_lock_worktree.assess_issue_lock_worktree(
worktree_path=dirty_repo["worktree"],
current_branch=BRANCH,
porcelain_status=porcelain,
base_equivalent=False,
)
assert assessment["block"] is True
assert any(
"tracked file edits exist before issue lock" in r
for r in assessment["reasons"]
)
# ── 15. Fixture matching #860 class with 7 fingerprint-pinned dirty paths ───
def test_issue_860_regression_fixture_spec():
spec = rebind.build_issue_860_regression_fixture_spec()
assert spec["claimant_identity"] == "jcwalker3"
assert spec["claimant_profile"] == "prgs-author"
assert spec["old_pid_alive"] is False
assert spec["live_session_pointer"] is None
assert spec["dirty_path_count"] == 7
assert len(spec["expected_dirty_paths"]) == 7
assert len(spec["expected_fingerprints"]) == 7
assert spec["expected_local_head"] == spec["expected_remote_head"]
for path in spec["expected_dirty_paths"]:
assert path in spec["expected_fingerprints"]
assert len(spec["expected_fingerprints"][path]) == 64
def test_dry_run_does_not_write(dirty_repo, lock_dir):
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
before = ils.read_lock_file(lock["lock_file_path"])
result = rebind.apply_dirty_same_claimant_session_rebind(
**_apply_kwargs(dirty_repo, lock, lock_dir, old_pid=old, dry_run=True)
)
assert result["success"], result
assert result["dry_run"] is True
after = ils.read_lock_file(lock["lock_file_path"])
assert after["session_pid"] == before["session_pid"]
assert not os.path.exists(os.path.join(lock_dir, f"session-{os.getpid()}.json"))
def test_provenance_source_is_sanctioned():
assert (
issue_lock_provenance.SOURCE_DIRTY_SAME_CLAIMANT_REBIND
in issue_lock_provenance.SANCTIONED_LOCK_SOURCES
)
assessment = issue_lock_provenance.assess_lock_file_for_create_pr(
{
"work_lease": {"operation_type": "author_issue_work"},
"lock_provenance": {
"source": issue_lock_provenance.SOURCE_DIRTY_SAME_CLAIMANT_REBIND,
"written_by_tool": issue_lock_provenance.SOURCE_DIRTY_SAME_CLAIMANT_REBIND,
"written_at": "2026-01-01T00:00:00Z",
},
}
)
assert assessment["proven"] is True
def test_permission_allowed_is_not_ownership_proof(dirty_repo, lock_dir):
"""permission_allowed=True must not bypass foreign claimant refusal."""
old = dead_pid()
lock = _make_lock(worktree=dirty_repo["worktree"], pid=old, lock_dir=lock_dir)
result = rebind.assess_dirty_same_claimant_session_rebind(
remote=REMOTE,
org=ORG,
repo=REPO,
issue_number=ISSUE,
branch_name=BRANCH,
worktree_path=dirty_repo["worktree"],
claimant_identity=IDENTITY,
claimant_profile=PROFILE,
old_pid=old,
expected_local_head=dirty_repo["local_head"],
expected_remote_head=dirty_repo["remote_head"],
expected_dirty_paths=dirty_repo["dirty_paths"],
expected_fingerprints=dirty_repo["fingerprints"],
existing_lock=lock,
current_identity="intruder",
current_profile=PROFILE,
role_kind="author",
current_pid=os.getpid(),
current_branch=BRANCH,
local_head=dirty_repo["local_head"],
remote_head=dirty_repo["remote_head"],
dirty_inventory=dirty_repo["inventory"],
permission_allowed=True,
repo_root=dirty_repo["root"],
)
assert not result["rebind_sanctioned"]
assert any("does not match active identity" in r for r in result["reasons"])
def test_content_fingerprint_stable(tmp_path):
p = tmp_path / "f.txt"
p.write_bytes(b"abc123")
a = rebind.content_fingerprint(str(p))
b = rebind.content_fingerprint(str(p))
assert a == b
assert len(a) == 64
+444
View File
@@ -0,0 +1,444 @@
"""Task heartbeat through the native MCP author path (#790 Slice A, AC-N6).
Assessor-level coverage is not sufficient here, and this project has already
paid for learning that: in review #499 on PR #791 the #760 renewal waiver was
computed correctly and then *discarded* at two later gates, so every real
renewal still failed while the unit suite stayed green. AC-N6 exists because of
that, and requires driving the real tools against a real git repository and a
real durable lock file, composing the gates in production order.
These tests therefore call ``gitea_lock_issue`` and
``gitea_heartbeat_issue_lock`` themselves and assert on what lands on disk,
never on an assessor's return value alone.
"""
from __future__ import annotations
import os
import subprocess
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from unittest.mock import patch
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from mutation_profile_fixture import shared_mutation_env # noqa: E402
import issue_lock_provenance # noqa: E402
import issue_lock_store # noqa: E402
import lease_policy # noqa: E402
import mcp_server # noqa: E402
ISSUE = 9791
BRANCH = f"fix/issue-{ISSUE}-heartbeat-mcp"
IDENTITY = "example-user"
PROFILE = "test-author-prgs"
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
def _ts(moment: datetime) -> str:
return (
moment.astimezone(timezone.utc)
.replace(microsecond=0)
.isoformat()
.replace("+00:00", "Z")
)
class _HeartbeatMcpBase(unittest.TestCase):
"""Real git repo plus a real durable lock, driven through the real tools."""
def setUp(self):
self.lock_dir = tempfile.TemporaryDirectory()
self.addCleanup(self.lock_dir.cleanup)
self.repo = tempfile.mkdtemp(prefix="issue790-mcp-")
self.addCleanup(lambda: subprocess.run(["rm", "-rf", self.repo], check=False))
self._init_worktree()
self.remotes = patch.dict(
mcp_server.REMOTES,
{"prgs": {"host": "gitea.prgs.cc", "org": ORG, "repo": REPO}},
)
self.remotes.start()
self.addCleanup(patch.stopall)
mcp_server._IDENTITY_CACHE.clear()
def _git(self, *args):
return subprocess.run(
["git", "-C", self.repo, *args], capture_output=True, text=True, check=True
)
def _init_worktree(self):
self._git("init", "-q", "-b", "master")
self._git("config", "user.email", "[email protected]")
self._git("config", "user.name", "Test")
with open(os.path.join(self.repo, "seed.txt"), "w") as fh:
fh.write("seed\n")
self._git("add", "seed.txt")
self._git("commit", "-q", "-m", "seed")
self.base_sha = self._git("rev-parse", "HEAD").stdout.strip()
# A fresh claim starts base-equivalent, which is the ordinary first-lock
# shape and exercises assess_issue_lock_worktree on its normal path.
self._git("checkout", "-q", "-b", BRANCH)
self.head_sha = self.base_sha
self.worktree = os.path.realpath(self.repo)
def _lock_path(self):
return issue_lock_store.lock_file_path(
remote="prgs",
org=ORG,
repo=REPO,
issue_number=ISSUE,
lock_dir=self.lock_dir.name,
)
def _tool_env(self):
env = shared_mutation_env(
PROFILE, include_example_repo=True, GITEA_ISSUE_LOCK_DIR=self.lock_dir.name
)
env["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
return env
def _git_state(self, *, porcelain="", base_equivalent=True):
return {
"current_branch": BRANCH,
"porcelain_status": porcelain,
"base_equivalent": base_equivalent,
"head_sha": self.head_sha,
"inspected_git_root": self.worktree,
"base_branch": "master",
}
def run_lock_issue(
self,
*,
branch_entries=None,
open_prs=None,
git_state=None,
identity=IDENTITY,
profile=PROFILE,
):
branch_entries = branch_entries if branch_entries is not None else []
open_prs = open_prs if open_prs is not None else []
git_state = git_state or self._git_state()
env = self._tool_env()
with patch(
"mcp_server.api_get_all", return_value=list(branch_entries)
), patch(
"mcp_server._list_open_pulls", return_value=list(open_prs)
), patch(
"mcp_server.get_auth_header", return_value="token x"
), patch(
"mcp_server._work_lease_claimant",
return_value={"username": identity, "profile": profile},
), patch(
"mcp_server.issue_lock_worktree.read_worktree_git_state",
return_value=git_state,
), patch(
"mcp_server.issue_duplicate_context_fetcher",
side_effect=lambda h, o, r, auth, issue_number: (
list(open_prs),
[b.get("name") for b in branch_entries if isinstance(b, dict)],
{"status": "not_claimed"},
),
), patch.dict(os.environ, env, clear=True):
os.environ["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
return mcp_server.gitea_lock_issue(
issue_number=ISSUE,
branch_name=BRANCH,
remote="prgs",
worktree_path=self.worktree,
)
def run_heartbeat(
self, *, task_session_id, identity=IDENTITY, profile=PROFILE, **kwargs
):
env = self._tool_env()
with patch(
"mcp_server._work_lease_claimant",
return_value={"username": identity, "profile": profile},
), patch("mcp_server.get_auth_header", return_value="token x"), patch.dict(
os.environ, env, clear=True
):
os.environ["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
return mcp_server.gitea_heartbeat_issue_lock(
issue_number=ISSUE,
branch_name=kwargs.pop("branch_name", BRANCH),
task_session_id=task_session_id,
remote="prgs",
worktree_path=kwargs.pop("worktree_path", self.worktree),
**kwargs,
)
def write_legacy_lock(self, *, hours_old: float = 3.0, ttl_hours: float = 4.0):
"""A durable lock in the shape the store wrote before this slice."""
now = datetime.now(timezone.utc)
claimant = {"username": IDENTITY, "profile": PROFILE}
created = now - timedelta(hours=hours_old)
record = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"remote": "prgs",
"org": ORG,
"repo": REPO,
"worktree_path": self.worktree,
"session_pid": os.getpid(),
"pid": os.getpid(),
"lock_generation": 1,
"work_lease": {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"pr_number": None,
"branch": BRANCH,
"worktree_path": self.worktree,
"claimant": claimant,
"created_at": _ts(created),
# The legacy signature: never advanced past creation.
"last_heartbeat_at": _ts(created),
"expires_at": _ts(created + timedelta(hours=ttl_hours)),
},
"lock_provenance": issue_lock_provenance.build_sanctioned_lock_provenance(
tool="gitea_lock_issue", claimant=claimant
),
}
path = self._lock_path()
record["lock_file_path"] = path
issue_lock_store.save_lock_file(path, record)
return record
class TestLockIssueMintsTheLifecycle(_HeartbeatMcpBase):
"""Durable lock creation and read-back through the real tool."""
def test_native_lock_writes_the_marker_and_a_task_session_id(self):
result = self.run_lock_issue()
self.assertTrue(result["success"], result)
written = issue_lock_store.read_lock_file(result["lock_file_path"])
lease = written["work_lease"]
self.assertEqual(
lease["lifecycle_version"], lease_policy.LIFECYCLE_HEARTBEAT_V1
)
self.assertTrue(lease["task_session_id"])
self.assertFalse(issue_lock_store.is_legacy_lease(written))
# AC-N1: the ownership key is not the daemon pid, which is recorded
# separately as evidence.
self.assertNotIn(str(written["session_pid"]), lease["task_session_id"])
self.assertEqual(written["session_pid"], os.getpid())
def test_native_lease_uses_the_policy_window_not_four_hours(self):
result = self.run_lock_issue()
lease = result["work_lease"]
created = datetime.fromisoformat(lease["created_at"].replace("Z", "+00:00"))
expires = datetime.fromisoformat(lease["expires_at"].replace("Z", "+00:00"))
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK)
self.assertEqual(
(expires - created).total_seconds() / 60.0, policy.initial_ttl_minutes
)
def test_freshness_of_a_new_native_lock_is_live(self):
result = self.run_lock_issue()
self.assertEqual(
result["lock_freshness"]["status"], issue_lock_store.STATUS_LIVE
)
self.assertTrue(result["lock_freshness"]["live"])
class TestHeartbeatThroughTheTool(_HeartbeatMcpBase):
def _lock_and_session(self):
result = self.run_lock_issue()
self.assertTrue(result["success"], result)
return result, result["work_lease"]["task_session_id"]
def test_heartbeat_slides_the_lease_and_advances_the_generation(self):
locked, session = self._lock_and_session()
before = issue_lock_store.read_lock_file(locked["lock_file_path"])
beat = self.run_heartbeat(task_session_id=session)
self.assertTrue(beat["success"], beat)
self.assertEqual(beat["operation"], "heartbeat")
after = issue_lock_store.read_lock_file(locked["lock_file_path"])
self.assertGreater(
issue_lock_store.lock_generation(after),
issue_lock_store.lock_generation(before),
)
self.assertGreaterEqual(
after["work_lease"]["expires_at"], before["work_lease"]["expires_at"]
)
self.assertEqual(after["work_lease"]["heartbeat_count"], 2)
def test_heartbeat_evidence_survives_the_downstream_mutation_gate(self):
"""The #499 F2 lesson, applied.
A sanction that is computed and then discarded downstream is worthless.
After a heartbeat the lock must still satisfy the gate every author
mutation runs through.
"""
locked, session = self._lock_and_session()
self.run_heartbeat(task_session_id=session)
written = issue_lock_store.read_lock_file(locked["lock_file_path"])
verdict = issue_lock_store.verify_lock_for_mutation(
written,
issue_number=ISSUE,
branch_name=BRANCH,
worktree_path=self.worktree,
)
self.assertTrue(verdict["proven"], verdict)
self.assertFalse(verdict["block"])
def _duplicate_gate(self, *, open_prs, branches):
env = self._tool_env()
with patch("mcp_server.get_auth_header", return_value="token x"), patch(
"mcp_server.issue_duplicate_context_fetcher",
side_effect=lambda h, o, r, auth, issue_number: (
list(open_prs),
list(branches),
{"status": "not_claimed"},
),
), patch.dict(os.environ, env, clear=True):
os.environ["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
return mcp_server.gitea_assess_work_issue_duplicate(
issue_number=ISSUE, branch_name=BRANCH, remote="prgs"
)
def test_heartbeat_does_not_change_the_duplicate_gate_verdict(self):
"""The gate must be invariant under heartbeating.
The point is not that the gate passes with a linked open PR at the
lock phase it correctly blocks (#400), heartbeat or not. The property
that matters is that sliding a lease neither loosens the gate nor
corrupts the lock state it reads: the verdict before and after a
heartbeat must be identical, for both the clear and the blocking shape.
"""
_, session = self._lock_and_session()
linked = [{"number": 4242, "head": {"ref": BRANCH, "sha": self.head_sha}}]
clear_before = self._duplicate_gate(open_prs=[], branches=[])
blocked_before = self._duplicate_gate(open_prs=linked, branches=[BRANCH])
self.assertTrue(self.run_heartbeat(task_session_id=session)["success"])
clear_after = self._duplicate_gate(open_prs=[], branches=[])
blocked_after = self._duplicate_gate(open_prs=linked, branches=[BRANCH])
self.assertEqual(clear_before["outcome"], clear_after["outcome"])
self.assertFalse(clear_after["block"])
self.assertEqual(blocked_before["outcome"], blocked_after["outcome"])
self.assertTrue(blocked_after["block"])
self.assertEqual(blocked_after["linked_open_pr"], 4242)
def test_foreign_session_id_is_refused_through_the_tool(self):
self._lock_and_session()
beat = self.run_heartbeat(task_session_id="author_issue_work-ffffffffffffffff")
self.assertFalse(beat["success"])
self.assertIn("task_session_id does not match", " ".join(beat["reasons"]))
def test_stale_generation_is_refused_through_the_tool(self):
locked, session = self._lock_and_session()
current = issue_lock_store.lock_generation(
issue_lock_store.read_lock_file(locked["lock_file_path"])
)
beat = self.run_heartbeat(
task_session_id=session, expected_generation=current + 5
)
self.assertFalse(beat["success"])
self.assertIn("generation changed", beat["reasons"][0])
def test_foreign_claimant_is_refused_through_the_tool(self):
_, session = self._lock_and_session()
beat = self.run_heartbeat(task_session_id=session, identity="someone-else")
self.assertFalse(beat["success"])
def test_heartbeat_cannot_acquire_a_missing_lock(self):
beat = self.run_heartbeat(task_session_id="author_issue_work-000000000000")
self.assertFalse(beat["success"])
self.assertIn("no durable lock", beat["reasons"][0])
def test_alive_pid_alone_does_not_keep_a_lease_live_through_the_tool(self):
"""PID-only refusal, end to end.
The recorded pid is this live process. The lock is aged past its grace
with no heartbeat, so the tool must refuse to slide it and the durable
record must classify as a missed heartbeat rather than as live.
"""
locked, session = self._lock_and_session()
record = issue_lock_store.read_lock_file(locked["lock_file_path"])
record["work_lease"]["last_heartbeat_at"] = _ts(
datetime.now(timezone.utc) - timedelta(minutes=30)
)
record["work_lease"]["expires_at"] = _ts(
datetime.now(timezone.utc) + timedelta(hours=2)
)
issue_lock_store.save_lock_file(locked["lock_file_path"], record)
self.assertTrue(issue_lock_store.is_process_alive(record["session_pid"]))
fresh = issue_lock_store.assess_lock_freshness(record)
self.assertEqual(
fresh["status"], issue_lock_store.STATUS_STALE_MISSED_HEARTBEAT
)
self.assertTrue(fresh["pid_alive"])
beat = self.run_heartbeat(task_session_id=session)
self.assertFalse(beat["success"])
self.assertIn("reclaimed", " ".join(beat["reasons"]))
class TestLegacyLocksThroughTheTool(_HeartbeatMcpBase):
"""AC-N8 end to end: protected on deployment, and rebindable."""
def test_legacy_lock_stays_protected_after_deployment(self):
record = self.write_legacy_lock(hours_old=3.0, ttl_hours=4.0)
fresh = issue_lock_store.assess_lock_freshness(record)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_LIVE)
self.assertTrue(fresh["legacy_lease"])
self.assertTrue(fresh["legacy_expiry_preserved"])
# It had never heartbeated, so under the new grace alone it would be
# long gone; the preserved absolute expiry is what protects it.
self.assertEqual(
record["work_lease"]["created_at"],
record["work_lease"]["last_heartbeat_at"],
)
def test_tool_rebinds_a_legacy_lock_and_mints_a_first_heartbeat(self):
self.write_legacy_lock(hours_old=3.0, ttl_hours=4.0)
result = self.run_heartbeat(task_session_id=None)
self.assertTrue(result["success"], result)
self.assertEqual(result["operation"], "legacy_rebind")
self.assertTrue(result["task_session_id"])
written = issue_lock_store.read_lock_file(self._lock_path())
lease = written["work_lease"]
self.assertEqual(
lease["lifecycle_version"], lease_policy.LIFECYCLE_HEARTBEAT_V1
)
self.assertEqual(lease["heartbeat_count"], 1)
self.assertNotEqual(
lease["created_at"],
written["legacy_rebind"]["legacy_origin"]["created_at"],
)
self.assertFalse(issue_lock_store.is_legacy_lease(written))
def test_rebound_lock_then_heartbeats_through_the_tool(self):
self.write_legacy_lock(hours_old=3.0, ttl_hours=4.0)
rebound = self.run_heartbeat(task_session_id=None)
beat = self.run_heartbeat(task_session_id=rebound["task_session_id"])
self.assertTrue(beat["success"], beat)
self.assertEqual(beat["operation"], "heartbeat")
self.assertEqual(beat["heartbeat_count"], 2)
def test_rebind_refuses_a_foreign_owner_through_the_tool(self):
self.write_legacy_lock(hours_old=3.0, ttl_hours=4.0)
result = self.run_heartbeat(task_session_id=None, identity="someone-else")
self.assertFalse(result["success"])
self.assertEqual(result["operation"], "legacy_rebind")
if __name__ == "__main__":
unittest.main()
+594
View File
@@ -0,0 +1,594 @@
"""Central lease policy and load-bearing heartbeat freshness (#790 Slice A).
Before this slice, ``issue_lock_store.assess_lock_freshness`` parsed
``last_heartbeat_at`` and then never consulted it: liveness was decided by an
absolute four-hour ``expires_at`` and by PID liveness. Because the recorded PID
is the long-lived MCP daemon rather than the authoring task, an abandoned claim
stayed "live" for the full four hours, and a claim whose work had already landed
blocked reconciliation for just as long (Issue #787 / PR #789, and again Issue
#760 / PR #791).
These tests pin the corrected semantics, including the two asymmetries that are
easy to lose in a refactor:
* an **alive** PID must never make anything live (AC-N2), while
* a **dead** PID must still mark a lease stale, because #753 dead-session
recovery keys on exactly that classification.
Durable-state helpers here write real lock files through the real flock path;
they are not mocks of the store.
"""
from __future__ import annotations
import os
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from unittest.mock import patch
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
import issue_lock_store # noqa: E402
import lease_policy # noqa: E402
import pr_work_lease # noqa: E402
import reviewer_pr_lease # noqa: E402
ISSUE = 9790
BRANCH = f"fix/issue-{ISSUE}-heartbeat"
IDENTITY = "example-user"
PROFILE = "test-author-prgs"
ORG = "Example-Org"
REPO = "Example-Repo"
REMOTE = "prgs"
DEAD_PID = 2**22 # far above any live pid on a test host
def _ts(moment: datetime) -> str:
return (
moment.astimezone(timezone.utc)
.replace(microsecond=0)
.isoformat()
.replace("+00:00", "Z")
)
class _LockFixture(unittest.TestCase):
def setUp(self):
self.lock_dir = tempfile.TemporaryDirectory()
self.addCleanup(self.lock_dir.cleanup)
self.now = datetime.now(timezone.utc)
self.worktree = os.path.realpath(tempfile.mkdtemp(prefix="issue790-"))
self.addCleanup(patch.stopall)
def _path(self):
return issue_lock_store.lock_file_path(
remote=REMOTE,
org=ORG,
repo=REPO,
issue_number=ISSUE,
lock_dir=self.lock_dir.name,
)
def write_lock(
self,
*,
lifecycle: str | None = lease_policy.LIFECYCLE_HEARTBEAT_V1,
created_delta: timedelta = timedelta(minutes=1),
heartbeat_delta: timedelta = timedelta(minutes=1),
expires_delta: timedelta = timedelta(minutes=9),
pid: int | None = None,
task_session_id: str | None = "author_issue_work-aaaabbbbccccdddd",
generation: int = 1,
identity: str = IDENTITY,
profile: str = PROFILE,
branch: str = BRANCH,
worktree: str | None = None,
) -> dict:
"""Write a real durable lock and return the record.
Deltas are relative to ``self.now``; ``expires_delta`` is added, the
others subtracted, so "in the past" reads naturally at each call site.
"""
lease: dict = {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"pr_number": None,
"branch": branch,
"worktree_path": worktree or self.worktree,
"claimant": {"username": identity, "profile": profile},
"created_at": _ts(self.now - created_delta),
"last_heartbeat_at": _ts(self.now - heartbeat_delta),
"expires_at": _ts(self.now + expires_delta),
}
if lifecycle is not None:
lease["lifecycle_version"] = lifecycle
if task_session_id is not None:
lease["task_session_id"] = task_session_id
pid_value = os.getpid() if pid is None else pid
record = {
"issue_number": ISSUE,
"branch_name": branch,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"worktree_path": worktree or self.worktree,
"session_pid": pid_value,
"pid": pid_value,
"lock_generation": generation,
"work_lease": lease,
}
path = self._path()
record["lock_file_path"] = path
issue_lock_store.save_lock_file(path, record)
return record
class TestPolicyIsTheSingleSource(unittest.TestCase):
"""AC-N7: one authoritative configuration source for every duration."""
def test_author_policy_carries_the_agreed_values(self):
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK)
self.assertEqual(policy.initial_ttl_minutes, 10.0)
self.assertEqual(policy.heartbeat_cadence_minutes, 2.0)
self.assertEqual(policy.stale_warning_minutes, 5.0)
self.assertEqual(policy.missed_heartbeat_grace_minutes, 10.0)
self.assertEqual(policy.absolute_cap_hours, 8.0)
self.assertEqual(policy.recovery_grace_minutes, 10.0)
self.assertEqual(policy.terminal_race_drain_minutes, 2.0)
self.assertTrue(policy.terminal_retirement_eligible)
self.assertTrue(policy.heartbeat_lifecycle_active)
def test_the_four_hour_author_ttl_literal_is_gone(self):
"""The duplicated literal AC-N7 exists to remove."""
self.assertFalse(hasattr(issue_lock_store, "WORK_LEASE_TTL_HOURS"))
import gitea_mcp_server
self.assertFalse(hasattr(gitea_mcp_server, "WORK_LEASE_TTL_HOURS"))
def test_declared_reviewer_values_match_the_module_still_using_them(self):
"""Slice A declares reviewer/merger numbers without rewiring them.
Recording a value in two places is only safe if drift is detectable, so
this asserts the declaration still equals the constants #747 owns. When
Slice C migrates those call sites, this test becomes the proof the
migration changed nothing.
"""
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_REVIEWER_PR)
self.assertEqual(
policy.initial_ttl_minutes, float(reviewer_pr_lease.LEASE_TTL_MINUTES)
)
self.assertEqual(
policy.stale_warning_minutes,
float(reviewer_pr_lease.STALE_WARNING_MINUTES),
)
self.assertFalse(policy.heartbeat_lifecycle_active)
def test_declared_conflict_fix_value_matches_its_module(self):
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_CONFLICT_FIX)
self.assertEqual(
policy.initial_ttl_minutes,
float(pr_work_lease.DEFAULT_CONFLICT_FIX_TTL_MINUTES),
)
self.assertFalse(policy.heartbeat_lifecycle_active)
def test_environment_override_applies(self):
var = lease_policy.env_var_name(
lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK, "initial_ttl_minutes"
)
with patch.dict(os.environ, {var: "7"}):
self.assertEqual(
lease_policy.policy_for(
lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK
).initial_ttl_minutes,
7.0,
)
def test_unusable_override_falls_back_instead_of_minting_a_zero_lease(self):
"""A typo must not make every claim instantly reclaimable."""
var = lease_policy.env_var_name(
lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK, "initial_ttl_minutes"
)
for bad in ("0", "-5", "not-a-number", " "):
with self.subTest(value=bad), patch.dict(os.environ, {var: bad}):
self.assertEqual(
lease_policy.policy_for(
lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK
).initial_ttl_minutes,
10.0,
)
def test_unknown_task_class_does_not_raise(self):
policy = lease_policy.policy_for("something-new")
self.assertEqual(policy.task_class, lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK)
class TestLifecycleDiscrimination(_LockFixture):
"""AC-N8: the marker, never a timestamp, decides legacy vs heartbeat."""
def test_missing_marker_reads_as_legacy(self):
record = self.write_lock(lifecycle=None)
self.assertTrue(issue_lock_store.is_legacy_lease(record))
self.assertEqual(
issue_lock_store.lease_lifecycle_version(record),
lease_policy.LIFECYCLE_LEGACY,
)
def test_marker_present_reads_as_heartbeat_lifecycle(self):
record = self.write_lock()
self.assertFalse(issue_lock_store.is_legacy_lease(record))
def test_equal_created_and_heartbeat_never_implies_a_fresh_heartbeat(self):
"""The exact inversion AC-N8 forbids.
A legacy lock has ``last_heartbeat_at == created_at`` forever because
nothing ever advanced it. Reading that equality as "recently
heartbeated" would classify every never-heartbeated lock as fresh.
"""
legacy = self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=3),
heartbeat_delta=timedelta(hours=3),
)
lease = legacy["work_lease"]
self.assertEqual(lease["created_at"], lease["last_heartbeat_at"])
self.assertTrue(issue_lock_store.is_legacy_lease(legacy))
# A brand-new heartbeat lease has them equal too, so the equality
# carries no information in either direction.
fresh = self.write_lock(
created_delta=timedelta(seconds=0), heartbeat_delta=timedelta(seconds=0)
)
self.assertEqual(
fresh["work_lease"]["created_at"],
fresh["work_lease"]["last_heartbeat_at"],
)
self.assertFalse(issue_lock_store.is_legacy_lease(fresh))
def test_minted_session_id_contains_no_pid(self):
"""AC-N1: the ownership key must not be derived from the daemon pid."""
minted = issue_lock_store.mint_task_session_id()
self.assertNotIn(str(os.getpid()), minted)
self.assertNotEqual(minted, issue_lock_store.mint_task_session_id())
class TestFreshnessIsHeartbeatDriven(_LockFixture):
"""AC-N2 and the new bands."""
def test_fresh_heartbeat_is_live(self):
record = self.write_lock(heartbeat_delta=timedelta(minutes=1))
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_LIVE)
self.assertTrue(fresh["live"])
self.assertFalse(fresh["heartbeat_warning"])
def test_heartbeat_past_warning_is_still_live_but_flagged(self):
record = self.write_lock(heartbeat_delta=timedelta(minutes=6))
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_LIVE)
self.assertTrue(fresh["heartbeat_warning"])
def test_missed_heartbeat_past_grace_is_classified_explicitly(self):
record = self.write_lock(
heartbeat_delta=timedelta(minutes=11),
expires_delta=timedelta(minutes=30),
)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(
fresh["status"], issue_lock_store.STATUS_STALE_MISSED_HEARTBEAT
)
self.assertFalse(fresh["live"])
self.assertTrue(fresh["stale"])
def test_alive_pid_never_establishes_freshness(self):
"""The defect in one assertion.
The recorded PID is this very process, so it is unambiguously alive
and the lease is still not live, because the task stopped heartbeating.
"""
record = self.write_lock(
pid=os.getpid(),
heartbeat_delta=timedelta(hours=4),
expires_delta=timedelta(hours=4),
)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertTrue(fresh["pid_alive"])
self.assertFalse(fresh["live"])
self.assertEqual(
fresh["status"], issue_lock_store.STATUS_STALE_MISSED_HEARTBEAT
)
def test_dead_pid_still_marks_stale_for_issue_753(self):
"""The opposite asymmetry: dead-PID corroboration is preserved."""
record = self.write_lock(pid=DEAD_PID, heartbeat_delta=timedelta(minutes=1))
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_STALE)
self.assertFalse(fresh["live"])
self.assertIn("not alive", fresh["reason"])
def test_absolute_cap_requires_readoption(self):
record = self.write_lock(
created_delta=timedelta(hours=9), heartbeat_delta=timedelta(minutes=1)
)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_STALE_ABSOLUTE_CAP)
self.assertIn("re-adoption", fresh["reason"])
def test_heartbeat_lifecycle_without_a_heartbeat_fails_closed(self):
record = self.write_lock()
del record["work_lease"]["last_heartbeat_at"]
issue_lock_store.save_lock_file(self._path(), record)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(
fresh["status"], issue_lock_store.STATUS_STALE_MISSED_HEARTBEAT
)
self.assertIn("fail closed", fresh["reason"])
def test_absent_lock(self):
fresh = issue_lock_store.assess_lock_freshness(None)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_ABSENT)
self.assertFalse(fresh["stale"])
class TestLegacyLocksStayProtected(_LockFixture):
"""AC-N8: deployment must not retroactively shorten an existing claim."""
def test_legacy_lock_with_a_stale_heartbeat_remains_live(self):
"""The deployment-safety case.
A four-hour legacy lease minted three hours ago has not heartbeated
once. Under the new grace it would be long gone; under its preserved
absolute expiry it is still live, and must stay that way.
"""
record = self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=3),
heartbeat_delta=timedelta(hours=3),
expires_delta=timedelta(hours=1),
)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_LIVE)
self.assertTrue(fresh["live"])
self.assertTrue(fresh["legacy_lease"])
self.assertTrue(fresh["legacy_expiry_preserved"])
def test_legacy_lock_past_its_absolute_expiry_is_expired_as_before(self):
record = self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=5),
heartbeat_delta=timedelta(hours=5),
expires_delta=timedelta(hours=-1),
)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_EXPIRED)
def test_legacy_lock_is_never_reclaimed_by_the_heartbeat_band(self):
record = self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=3),
heartbeat_delta=timedelta(hours=3),
expires_delta=timedelta(hours=1),
)
reclaim = issue_lock_store.assess_expired_lock_reclaim(record, now=self.now)
self.assertFalse(reclaim["reclaim_allowed"])
class TestReclaimAfterMissedHeartbeat(_LockFixture):
def test_missed_heartbeat_makes_ownership_reclaimable(self):
record = self.write_lock(
pid=os.getpid(),
heartbeat_delta=timedelta(minutes=15),
expires_delta=timedelta(hours=3),
)
reclaim = issue_lock_store.assess_expired_lock_reclaim(record, now=self.now)
self.assertTrue(reclaim["reclaim_allowed"])
self.assertIn("stale_missed_heartbeat", reclaim["reasons"][0])
def test_live_lease_is_never_reclaimable(self):
record = self.write_lock(heartbeat_delta=timedelta(minutes=1))
reclaim = issue_lock_store.assess_expired_lock_reclaim(record, now=self.now)
self.assertFalse(reclaim["reclaim_allowed"])
def test_dead_pid_reclaim_path_is_unchanged(self):
"""#753 must keep working through its original conditions."""
record = self.write_lock(pid=DEAD_PID, heartbeat_delta=timedelta(minutes=1))
reclaim = issue_lock_store.assess_expired_lock_reclaim(record, now=self.now)
self.assertTrue(reclaim["reclaim_allowed"])
self.assertTrue(reclaim["owner_pid_dead"])
class TestHeartbeatWriter(_LockFixture):
"""A4: flock + CAS + exact verification, and no revival path."""
def _heartbeat(self, **kwargs):
params = {
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": self.worktree,
"identity": IDENTITY,
"profile": PROFILE,
"task_session_id": "author_issue_work-aaaabbbbccccdddd",
"lock_dir": self.lock_dir.name,
"now": self.now,
}
params.update(kwargs)
return issue_lock_store.heartbeat_session_lock(**params)
def test_heartbeat_slides_expiry_and_advances_generation(self):
self.write_lock(heartbeat_delta=timedelta(minutes=4), generation=5)
result = self._heartbeat()
self.assertTrue(result["success"], result)
self.assertEqual(result["prior_generation"], 5)
self.assertEqual(result["lock_generation"], 6)
self.assertEqual(result["heartbeat_count"], 1)
self.assertEqual(result["last_heartbeat_at"], _ts(self.now))
self.assertEqual(result["expires_at"], _ts(self.now + timedelta(minutes=10)))
self.assertTrue(result["freshness"]["live"])
def test_heartbeat_is_durable_and_repeatable(self):
self.write_lock(heartbeat_delta=timedelta(minutes=4))
self._heartbeat()
second = self._heartbeat(now=self.now + timedelta(minutes=1))
self.assertTrue(second["success"], second)
self.assertEqual(second["heartbeat_count"], 2)
written = issue_lock_store.read_lock_file(self._path())
self.assertEqual(written["work_lease"]["heartbeat_count"], 2)
def test_stale_generation_is_refused(self):
self.write_lock(generation=5)
result = self._heartbeat(expected_generation=4)
self.assertFalse(result["success"])
self.assertIn("generation changed", result["reasons"][0])
def test_foreign_session_is_refused(self):
self.write_lock()
result = self._heartbeat(task_session_id="author_issue_work-ffffffffffffffff")
self.assertFalse(result["success"])
self.assertIn("task_session_id does not match", " ".join(result["reasons"]))
def test_missing_session_id_is_refused(self):
self.write_lock()
result = self._heartbeat(task_session_id="")
self.assertFalse(result["success"])
def test_foreign_claimant_is_refused(self):
self.write_lock()
for field, value in (
("identity", "someone-else"),
("profile", "other-profile"),
):
with self.subTest(field=field):
result = self._heartbeat(**{field: value})
self.assertFalse(result["success"])
def test_branch_and_worktree_mismatch_are_refused(self):
self.write_lock()
wrong_branch = self._heartbeat(branch_name=f"fix/issue-{ISSUE}-other")
self.assertFalse(wrong_branch["success"])
wrong_worktree = self._heartbeat(worktree_path="/tmp/not-the-worktree")
self.assertFalse(wrong_worktree["success"])
def test_lapsed_lease_cannot_be_heartbeated_back_to_life(self):
"""No revival path (A4).
A session that stopped proving liveness must reclaim under a fresh
generation, not restore ownership retroactively.
"""
self.write_lock(
heartbeat_delta=timedelta(minutes=30), expires_delta=timedelta(hours=1)
)
result = self._heartbeat()
self.assertFalse(result["success"])
self.assertIn("reclaimed", " ".join(result["reasons"]))
def test_absent_lock_cannot_be_created_by_heartbeat(self):
result = self._heartbeat()
self.assertFalse(result["success"])
self.assertIn("no durable lock", result["reasons"][0])
def test_legacy_lock_is_refused_until_rebound(self):
self.write_lock(lifecycle=None)
result = self._heartbeat()
self.assertFalse(result["success"])
self.assertTrue(result["legacy_lease"])
self.assertIn("rebound", " ".join(result["reasons"]))
class TestLegacyRebind(_LockFixture):
"""AC-N8 exit route: canonical exact-owner rebinding."""
def _rebind(self, **kwargs):
params = {
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": self.worktree,
"identity": IDENTITY,
"profile": PROFILE,
"lock_dir": self.lock_dir.name,
"now": self.now,
}
params.update(kwargs)
return issue_lock_store.rebind_legacy_lock(**params)
def test_rebind_mints_a_session_and_a_genuine_first_heartbeat(self):
self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=3),
heartbeat_delta=timedelta(hours=3),
expires_delta=timedelta(hours=1),
generation=2,
)
result = self._rebind()
self.assertTrue(result["success"], result)
self.assertTrue(result["task_session_id"])
self.assertEqual(result["lock_generation"], 3)
written = issue_lock_store.read_lock_file(self._path())
lease = written["work_lease"]
self.assertEqual(
lease["lifecycle_version"], lease_policy.LIFECYCLE_HEARTBEAT_V1
)
self.assertEqual(lease["last_heartbeat_at"], _ts(self.now))
self.assertEqual(lease["expires_at"], _ts(self.now + timedelta(minutes=10)))
self.assertFalse(issue_lock_store.is_legacy_lease(written))
# The original claim is preserved for audit rather than overwritten.
origin = written["legacy_rebind"]["legacy_origin"]
self.assertTrue(origin["created_at"])
self.assertEqual(origin["lifecycle"], lease_policy.LIFECYCLE_LEGACY)
def test_rebound_lock_can_then_heartbeat(self):
self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=3),
heartbeat_delta=timedelta(hours=3),
expires_delta=timedelta(hours=1),
)
rebound = self._rebind()
beat = issue_lock_store.heartbeat_session_lock(
remote=REMOTE,
org=ORG,
repo=REPO,
issue_number=ISSUE,
branch_name=BRANCH,
worktree_path=self.worktree,
identity=IDENTITY,
profile=PROFILE,
task_session_id=rebound["task_session_id"],
lock_dir=self.lock_dir.name,
now=self.now + timedelta(minutes=1),
)
self.assertTrue(beat["success"], beat)
def test_rebind_refuses_a_foreign_owner(self):
self.write_lock(lifecycle=None, expires_delta=timedelta(hours=1))
result = self._rebind(identity="someone-else")
self.assertFalse(result["success"])
def test_rebind_refuses_a_lock_already_on_the_lifecycle(self):
self.write_lock()
result = self._rebind()
self.assertFalse(result["success"])
self.assertFalse(result["legacy_lease"])
def test_rebind_is_not_a_recovery_path_for_a_lapsed_legacy_lease(self):
"""An expired legacy lease belongs to #760 renewal or #601 reclaim."""
self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=5),
heartbeat_delta=timedelta(hours=5),
expires_delta=timedelta(hours=-1),
)
result = self._rebind()
self.assertFalse(result["success"])
self.assertIn("not a recovery path", " ".join(result["reasons"]))
if __name__ == "__main__":
unittest.main()
@@ -1,244 +0,0 @@
"""#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()
@@ -1,551 +0,0 @@
"""Merged-PR awareness for the worktree cleanup audit (#858).
Before #858 an ``issue_work`` worktree could never leave ``active_issue_work``:
the audit had no PR linkage at all (``pr_number`` was structurally ``None``)
and its only route to ``clean_stale_removable`` was a TTL derived from a
``last_used_at`` that nothing ever populated. A merged, clean, unprotected
worktree was therefore reported as active work forever, disagreeing with the
PR-scoped reconciler.
These tests use fabricated temporary repositories and synthetic PR records
only. Nothing here removes a worktree or deletes a branch.
"""
import os
import subprocess
import sys
import tempfile
import unittest
from unittest.mock import patch
sys.path.insert(0, str(__import__("pathlib").Path(__file__).resolve().parent.parent))
import merged_cleanup_reconcile as mcr # noqa: E402
import worktree_cleanup_audit as wca # noqa: E402
MERGED_BRANCH = "feat/issue-777-timeline"
MERGED_PATH = "/repo/branches/issue-777-timeline"
HEAD_SHA = "a" * 40
def _pr(number, branch, *, merged=True, sha=HEAD_SHA, state=None):
"""Synthetic Gitea PR payload."""
return {
"number": number,
"head": {"ref": branch, "sha": sha},
"merged_at": "2026-07-24T01:00:00Z" if merged else None,
"state": state or ("closed" if merged else "open"),
}
def _porcelain(*entries):
out = []
for path, branch, sha in entries:
out.append(f"worktree {path}")
out.append(f"HEAD {sha}")
if branch is None:
out.append("detached")
else:
out.append(f"branch refs/heads/{branch}")
out.append("")
return "\n".join(out)
class _AuditHarness(unittest.TestCase):
"""Runs audit_branches_directory over a fabricated worktree listing."""
PORCELAIN = _porcelain(
("/repo", "master", "f" * 40),
(MERGED_PATH, MERGED_BRANCH, HEAD_SHA),
)
def run_audit(self, *, dirty_paths=(), contained=True, **kwargs):
def fake_dirty(path):
if path in dirty_paths:
return {"exists": True, "dirty": True, "dirty_files": [" M x.py"]}
return {"exists": True, "dirty": False, "dirty_files": []}
with patch.object(
wca, "list_worktrees",
return_value=wca.parse_worktree_porcelain(self.PORCELAIN),
), patch.object(
wca, "read_worktree_dirty", side_effect=fake_dirty
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=contained
):
report = wca.audit_branches_directory("/repo", **kwargs)
return {wt["path"]: wt for wt in report["worktrees"]}, report
def merged_audit(self, **kwargs):
kwargs.setdefault("pr_index", wca.build_pr_index([_pr(849, MERGED_BRANCH)]))
kwargs.setdefault("master_ref", "prgs/master")
return self.run_audit(**kwargs)
class TestMergedWorktreeBecomesRemovable(_AuditHarness):
def test_clean_merged_issue_worktree_is_linked_and_removable(self):
by_path, report = self.merged_audit()
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_CLEAN_STALE_REMOVABLE)
self.assertTrue(entry["removable"])
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_MERGED)
self.assertEqual(entry["merged_pr_cleanup"]["block_reasons"], [])
self.assertIn(MERGED_PATH, [c["path"] for c in report["removable_candidates"]])
def test_pr_number_populated_from_authoritative_linkage(self):
by_path, _ = self.merged_audit()
self.assertEqual(by_path[MERGED_PATH]["pr_number"], 849)
def test_regression_without_pr_evidence_stays_active_issue_work(self):
"""The pre-#858 behaviour, still correct when no PR state is supplied."""
by_path, _ = self.run_audit()
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
self.assertIsNone(entry["pr_number"])
class TestProtectiveSignalsSurvive(_AuditHarness):
def test_open_pr_worktree_is_not_removable(self):
index = wca.build_pr_index([_pr(900, MERGED_BRANCH, merged=False)])
by_path, _ = self.run_audit(
pr_index=index,
master_ref="prgs/master",
open_pr_branches={MERGED_BRANCH},
)
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_OPEN_PR)
self.assertFalse(entry["removable"])
# linkage still reports the owning PR, it just is not merge proof
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_OPEN)
self.assertEqual(entry["pr_number"], 900)
def test_dirty_tracked_worktree_is_not_removable(self):
by_path, _ = self.merged_audit(dirty_paths=(MERGED_PATH,))
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_DIRTY_LOCAL)
self.assertFalse(entry["removable"])
self.assertIn(
"worktree has uncommitted changes",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_untracked_only_worktree_is_not_removable(self):
"""``git status --porcelain`` reports untracked files as dirty too."""
def untracked(path):
if path == MERGED_PATH:
return {"exists": True, "dirty": True, "dirty_files": ["?? scratch.txt"]}
return {"exists": True, "dirty": False, "dirty_files": []}
with patch.object(
wca, "list_worktrees",
return_value=wca.parse_worktree_porcelain(self.PORCELAIN),
), patch.object(
wca, "read_worktree_dirty", side_effect=untracked
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory(
"/repo",
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
)
entry = {wt["path"]: wt for wt in report["worktrees"]}[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_DIRTY_LOCAL)
self.assertFalse(entry["removable"])
def test_active_lease_by_issue_number_is_protective(self):
by_path, _ = self.merged_audit(leased_issue_numbers={777})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_active_lease"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_active_lease_by_branch_is_protective(self):
by_path, _ = self.merged_audit(leased_branches={MERGED_BRANCH})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_active_lease"])
self.assertFalse(entry["removable"])
def test_active_issue_lock_is_protective(self):
by_path, _ = self.merged_audit(active_issue_branches={MERGED_BRANCH})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_active_issue_lock"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_live_session_worktree_is_protective(self):
by_path, _ = self.merged_audit(live_session_paths={MERGED_PATH})
entry = by_path[MERGED_PATH]
self.assertTrue(entry["has_live_session"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_head_not_contained_in_master_is_not_removable(self):
by_path, _ = self.merged_audit(contained=False)
entry = by_path[MERGED_PATH]
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
self.assertIn(
"worktree head is not contained in authoritative master "
"(unmerged commits remain)",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_unknown_containment_fails_closed(self):
by_path, _ = self.merged_audit(contained=None)
entry = by_path[MERGED_PATH]
self.assertFalse(entry["removable"])
self.assertIn(
"containment of the worktree head in master is unknown",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_missing_master_ref_fails_closed(self):
by_path, _ = self.run_audit(
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)])
)
self.assertFalse(by_path[MERGED_PATH]["removable"])
def test_unmerged_owning_pr_is_not_removable(self):
index = wca.build_pr_index([_pr(901, MERGED_BRANCH, merged=False)])
by_path, _ = self.run_audit(pr_index=index, master_ref="prgs/master")
entry = by_path[MERGED_PATH]
self.assertFalse(entry["removable"])
self.assertIn(
"owning PR #901 is not merged",
entry["merged_pr_cleanup"]["block_reasons"],
)
def test_control_checkout_is_never_removable(self):
by_path, _ = self.merged_audit()
control = by_path["/repo"]
self.assertTrue(control["is_protected"])
self.assertEqual(control["classification"], wca.CLASS_UNSAFE_UNKNOWN)
self.assertFalse(control["removable"])
def test_control_checkout_not_removable_even_if_linked_and_merged(self):
"""A merged PR on the control checkout must not unlock removal."""
porcelain = _porcelain(("/repo", MERGED_BRANCH, HEAD_SHA))
with patch.object(
wca, "list_worktrees", return_value=wca.parse_worktree_porcelain(porcelain)
), patch.object(
wca, "read_worktree_dirty",
return_value={"exists": True, "dirty": False, "dirty_files": []},
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory(
"/repo",
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
)
entry = report["worktrees"][0]
self.assertEqual(entry["classification"], wca.CLASS_UNSAFE_UNKNOWN)
self.assertFalse(entry["removable"])
class TestAmbiguousLinkageFailsClosed(_AuditHarness):
def test_competing_prs_on_one_branch_fail_closed(self):
index = wca.build_pr_index(
[_pr(849, MERGED_BRANCH), _pr(860, MERGED_BRANCH)]
)
by_path, _ = self.run_audit(pr_index=index, master_ref="prgs/master")
entry = by_path[MERGED_PATH]
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_AMBIGUOUS)
self.assertIsNone(entry["pr_number"])
self.assertEqual(entry["classification"], wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(entry["removable"])
def test_merged_plus_open_pr_on_one_branch_fails_closed(self):
index = wca.build_pr_index(
[_pr(849, MERGED_BRANCH), _pr(861, MERGED_BRANCH, merged=False)]
)
by_path, _ = self.run_audit(pr_index=index, master_ref="prgs/master")
entry = by_path[MERGED_PATH]
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_AMBIGUOUS)
self.assertFalse(entry["removable"])
def test_no_owning_pr_fails_closed(self):
by_path, _ = self.run_audit(
pr_index=wca.build_pr_index([_pr(849, "feat/other-branch")]),
master_ref="prgs/master",
)
entry = by_path[MERGED_PATH]
self.assertEqual(entry["merged_pr_linkage"]["status"], wca.LINKAGE_NONE)
self.assertFalse(entry["removable"])
def test_malformed_pr_records_are_dropped_not_guessed(self):
index = wca.build_pr_index(
[
{"number": None, "head": {"ref": MERGED_BRANCH}},
{"number": 5, "head": {}},
{"number": "not-an-int", "head": {"ref": MERGED_BRANCH}},
]
)
self.assertEqual(index, {})
self.assertEqual(
wca.resolve_owning_pr(branch=MERGED_BRANCH, pr_index=index)["status"],
wca.LINKAGE_NONE,
)
def test_detached_worktree_has_no_branch_linkage(self):
self.assertEqual(
wca.resolve_owning_pr(branch=None, pr_index={})["status"],
wca.LINKAGE_UNKNOWN,
)
class TestUnrelatedClassificationsUnchanged(unittest.TestCase):
"""Non-issue_work worktrees keep their pre-#858 classifications."""
PORCELAIN = _porcelain(
("/repo", "master", "f" * 40),
("/repo/branches/review-pr42", "review-pr42", "2" * 40),
("/repo/branches/baseline-master-x", "baseline-master-x", "3" * 40),
("/repo/branches/conflict-fix-pr50", "conflict-fix-pr50", "4" * 40),
("/repo/branches/review-pr99", None, "5" * 40),
)
def _audit(self, **kwargs):
with patch.object(
wca, "list_worktrees",
return_value=wca.parse_worktree_porcelain(self.PORCELAIN),
), patch.object(
wca, "read_worktree_dirty",
return_value={"exists": True, "dirty": False, "dirty_files": []},
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory("/repo", **kwargs)
return {wt["path"]: wt for wt in report["worktrees"]}
def test_classifications_identical_with_and_without_pr_evidence(self):
without = self._audit()
with_evidence = self._audit(
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
)
self.assertEqual(
{p: e["classification"] for p, e in without.items()},
{p: e["classification"] for p, e in with_evidence.items()},
)
def test_lease_on_issue_does_not_capture_similarly_named_scratch_trees(self):
"""A lease on issue 777 protects issue work, not baseline/review trees."""
porcelain = _porcelain(
("/repo/branches/baseline-master-issue-777", "baseline-issue-777", "7" * 40),
("/repo/branches/issue-777-timeline", MERGED_BRANCH, HEAD_SHA),
)
with patch.object(
wca, "list_worktrees", return_value=wca.parse_worktree_porcelain(porcelain)
), patch.object(
wca, "read_worktree_dirty",
return_value={"exists": True, "dirty": False, "dirty_files": []},
), patch.object(
wca, "git_worktree_list", return_value="(mocked)"
), patch.object(
wca, "is_head_ancestor_of_ref", return_value=True
):
report = wca.audit_branches_directory(
"/repo",
pr_index=wca.build_pr_index([_pr(849, MERGED_BRANCH)]),
master_ref="prgs/master",
leased_issue_numbers={777},
)
by_path = {wt["path"]: wt for wt in report["worktrees"]}
baseline = by_path["/repo/branches/baseline-master-issue-777"]
self.assertFalse(baseline["has_active_lease"])
self.assertEqual(baseline["classification"], wca.CLASS_CLEAN_STALE_REMOVABLE)
issue_work = by_path["/repo/branches/issue-777-timeline"]
self.assertTrue(issue_work["has_active_lease"])
self.assertFalse(issue_work["removable"])
def test_review_and_baseline_still_removable(self):
by_path = self._audit(
pr_index=wca.build_pr_index([]), master_ref="prgs/master"
)
self.assertEqual(
by_path["/repo/branches/review-pr42"]["classification"],
wca.CLASS_CLEAN_STALE_REMOVABLE,
)
self.assertEqual(
by_path["/repo/branches/baseline-master-x"]["classification"],
wca.CLASS_CLEAN_STALE_REMOVABLE,
)
self.assertEqual(
by_path["/repo/branches/review-pr99"]["classification"],
wca.CLASS_DETACHED_REVIEW_LEFTOVER,
)
def test_conflict_fix_ttl_behaviour_unchanged(self):
"""conflict_fix still needs only TTL expiry; #858 did not touch it."""
self.assertEqual(
wca.classify_worktree(
workflow_type=wca.WORKFLOW_CONFLICT_FIX,
is_dirty=False,
ttl_expired=True,
),
wca.CLASS_CLEAN_STALE_REMOVABLE,
)
self.assertEqual(
wca.classify_worktree(
workflow_type=wca.WORKFLOW_CONFLICT_FIX,
is_dirty=False,
ttl_expired=False,
),
wca.CLASS_ACTIVE_ISSUE_WORK,
)
def test_issue_work_ttl_alone_no_longer_grants_removal(self):
"""Age is not landing proof: TTL alone must not reclaim issue work."""
self.assertEqual(
wca.classify_worktree(
workflow_type=wca.WORKFLOW_ISSUE_WORK,
is_dirty=False,
ttl_expired=True,
),
wca.CLASS_ACTIVE_ISSUE_WORK,
)
class TestAssessorPerformsNoDeletion(_AuditHarness):
def test_audit_never_removes_a_worktree(self):
with patch.object(wca, "remove_worktree") as removal:
self.merged_audit()
removal.assert_not_called()
def test_audit_shells_out_to_no_destructive_git_command(self):
seen = []
real_run = subprocess.run
def recording_run(cmd, *args, **kwargs):
seen.append(cmd)
return real_run(["true"], *args, **kwargs)
with patch.object(subprocess, "run", side_effect=recording_run):
wca.audit_branches_directory("/nonexistent-repo-for-audit")
joined = [" ".join(c) if isinstance(c, list) else str(c) for c in seen]
for cmd in joined:
self.assertNotIn("worktree remove", cmd)
self.assertNotIn("branch -D", cmd)
self.assertNotIn("push", cmd)
class TestAgreementWithPrScopedReconciler(unittest.TestCase):
"""The audit and merged_cleanup_reconcile must agree on identical input.
Uses a real throwaway git repository so containment is computed by git
rather than asserted. Nothing outside the temporary directory is touched.
"""
def _git(self, *args):
subprocess.run(
["git", "-C", self.root, *args],
check=True,
capture_output=True,
text=True,
)
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.root = os.path.realpath(self._tmp.name)
self._git("init", "-b", "master", ".")
self._git("config", "user.email", "[email protected]")
self._git("config", "user.name", "Test")
with open(os.path.join(self.root, "seed.txt"), "w") as fh:
fh.write("seed\n")
self._git("add", "seed.txt")
self._git("commit", "-m", "seed")
self.branch = "feat/issue-777-timeline"
self._git("checkout", "-b", self.branch)
with open(os.path.join(self.root, "feature.txt"), "w") as fh:
fh.write("feature\n")
self._git("add", "feature.txt")
self._git("commit", "-m", "feature")
self.head_sha = subprocess.run(
["git", "-C", self.root, "rev-parse", "HEAD"],
capture_output=True, text=True, check=True,
).stdout.strip()
self._git("checkout", "master")
self._git("merge", "--no-ff", "-m", "merge feature", self.branch)
self.worktree = os.path.join(self.root, "branches", "issue-777-timeline")
self._git("worktree", "add", self.worktree, self.branch)
def tearDown(self):
self._tmp.cleanup()
def _pr_index(self):
return wca.build_pr_index(
[
{
"number": 849,
"head": {"ref": self.branch, "sha": self.head_sha},
"merged_at": "2026-07-24T01:00:00Z",
}
]
)
def _audit_entry(self):
report = wca.audit_branches_directory(
self.root, pr_index=self._pr_index(), master_ref="master"
)
return next(wt for wt in report["worktrees"] if wt["path"] == self.worktree)
def _reconciler_entry(self):
return mcr.assess_local_worktree_cleanup(
pr_number=849,
head_branch=self.branch,
merged=True,
worktree_state=mcr.resolve_cleanup_worktree_state(
project_root=self.root,
head_branch=self.branch,
issue_number=777,
pr_head_sha=self.head_sha,
target_ref="master",
),
active_lock=False,
)
def test_both_assessors_agree_the_worktree_is_safe(self):
audit_entry = self._audit_entry()
reconciler = self._reconciler_entry()
self.assertTrue(reconciler["safe_to_remove_worktree"], reconciler)
self.assertTrue(audit_entry["removable"], audit_entry)
self.assertEqual(audit_entry["pr_number"], reconciler["pr_number"])
self.assertEqual(audit_entry["merged_pr_cleanup"]["block_reasons"], [])
self.assertEqual(reconciler["block_reasons"], [])
def test_both_assessors_agree_a_dirty_worktree_is_unsafe(self):
with open(os.path.join(self.worktree, "feature.txt"), "a") as fh:
fh.write("local edit\n")
audit_entry = self._audit_entry()
reconciler = self._reconciler_entry()
self.assertFalse(audit_entry["removable"])
self.assertFalse(reconciler["safe_to_remove_worktree"])
def test_worktree_still_present_after_audit(self):
self._audit_entry()
self.assertTrue(os.path.isdir(self.worktree))
if __name__ == "__main__":
unittest.main()
+2 -24
View File
@@ -134,35 +134,13 @@ class TestClassification(unittest.TestCase):
self.assertEqual(cls, wca.CLASS_ACTIVE_OPEN_PR)
self.assertFalse(wca.is_removable(cls))
def test_stale_clean_issue_worktree_needs_merged_pr_proof(self):
# Scenario 5 (#858): age is not proof that the branch landed, so a
# TTL-expired issue worktree stays active work. Only authoritative
# merged-PR evidence makes it removable, which is what keeps a
# worktree holding unmerged commits from being reclaimed by age.
def test_stale_clean_issue_worktree_removable(self):
# Scenario 5: clean issue worktree, TTL expired, no lock -> removable.
cls = wca.classify_worktree(
workflow_type=wca.WORKFLOW_ISSUE_WORK,
is_dirty=False,
ttl_expired=True,
)
self.assertEqual(cls, wca.CLASS_ACTIVE_ISSUE_WORK)
self.assertFalse(wca.is_removable(cls))
cls = wca.classify_worktree(
workflow_type=wca.WORKFLOW_ISSUE_WORK,
is_dirty=False,
ttl_expired=True,
merged_pr_cleanup={"proven": True},
)
self.assertEqual(cls, wca.CLASS_CLEAN_STALE_REMOVABLE)
self.assertTrue(wca.is_removable(cls))
def test_stale_clean_conflict_fix_worktree_removable(self):
# conflict_fix keeps the original TTL rule; #858 changed issue work only.
cls = wca.classify_worktree(
workflow_type=wca.WORKFLOW_CONFLICT_FIX,
is_dirty=False,
ttl_expired=True,
)
self.assertEqual(cls, wca.CLASS_CLEAN_STALE_REMOVABLE)
self.assertTrue(wca.is_removable(cls))
+4 -271
View File
@@ -34,11 +34,7 @@ import subprocess
from datetime import datetime, timezone
from typing import Any
from merged_cleanup_reconcile import (
branch_worktree_folder,
is_head_ancestor_of_ref,
read_local_worktree_state,
)
from merged_cleanup_reconcile import branch_worktree_folder, read_local_worktree_state
from reviewer_worktree import parse_dirty_tracked_files, REVIEW_WORKTREE_RE
PROTECTED_BRANCHES = frozenset({"master", "main", "dev"})
@@ -71,14 +67,6 @@ REMOVABLE_CLASSES = frozenset(
{CLASS_CLEAN_STALE_REMOVABLE, CLASS_DETACHED_REVIEW_LEFTOVER}
)
# Merged-PR linkage outcomes for issue worktrees (#858). Only ``LINKAGE_MERGED``
# is ownership proof; every other outcome leaves the worktree protected.
LINKAGE_MERGED = "merged_pr"
LINKAGE_OPEN = "open_pr"
LINKAGE_NONE = "no_owning_pr"
LINKAGE_AMBIGUOUS = "ambiguous"
LINKAGE_UNKNOWN = "unknown"
_ISSUE_REF_RE = re.compile(r"issue-(\d+)", re.IGNORECASE)
_ISSUE_BRANCH_PREFIXES = ("feat/", "fix/", "docs/", "chore/")
@@ -181,186 +169,6 @@ def is_ttl_expired(
return (now_dt - last).total_seconds() > ttl_hours * 3600.0
def build_pr_index(prs: list[dict[str, Any]] | None) -> dict[str, list[dict[str, Any]]]:
"""Index PR records by head branch for deterministic worktree linkage (#858).
Accepts Gitea PR payloads (``head`` as a dict) and pre-flattened records
(``head_branch``/``head_sha``). Records without a usable head branch or
number are dropped rather than guessed at, so a branch is only ever linked
to a PR the caller actually proved.
"""
index: dict[str, list[dict[str, Any]]] = {}
for pr in prs or []:
head = pr.get("head")
if isinstance(head, dict):
head_branch = head.get("ref")
head_sha = head.get("sha")
else:
head_branch = pr.get("head_branch") or (head if isinstance(head, str) else None)
head_sha = pr.get("head_sha")
number = pr.get("number")
if not head_branch or number is None:
continue
try:
pr_number = int(number)
except (TypeError, ValueError):
continue
index.setdefault(str(head_branch).strip(), []).append(
{
"pr_number": pr_number,
"head_branch": str(head_branch).strip(),
"head_sha": head_sha,
"merged": bool(pr.get("merged") or pr.get("merged_at")),
"state": pr.get("state"),
}
)
return index
def resolve_owning_pr(
*,
branch: str | None,
pr_index: dict[str, list[dict[str, Any]]] | None,
) -> dict[str, Any]:
"""Resolve the single PR that owns ``branch``, failing closed when unclear.
Ownership is only ``LINKAGE_MERGED`` when exactly one PR claims the branch
and that PR is merged. Several distinct PRs on one branch is a competing
claim (``LINKAGE_AMBIGUOUS``), and a still-open owner is reported as
``LINKAGE_OPEN`` both keep the worktree protected while still exposing
the PR number the audit resolved.
"""
if pr_index is None:
return {
"status": LINKAGE_UNKNOWN,
"pr_number": None,
"candidate_pr_numbers": [],
"reasons": ["live PR state was not supplied; ownership unproven"],
}
branch_name = (branch or "").strip()
if not branch_name:
return {
"status": LINKAGE_UNKNOWN,
"pr_number": None,
"candidate_pr_numbers": [],
"reasons": ["worktree has no attached branch; ownership unproven"],
}
candidates = list(pr_index.get(branch_name) or [])
numbers = sorted({c["pr_number"] for c in candidates})
if not candidates:
return {
"status": LINKAGE_NONE,
"pr_number": None,
"candidate_pr_numbers": [],
"reasons": [f"no PR claims branch '{branch_name}'"],
}
if len(numbers) > 1:
return {
"status": LINKAGE_AMBIGUOUS,
"pr_number": None,
"candidate_pr_numbers": numbers,
"reasons": [
f"branch '{branch_name}' is claimed by competing PRs {numbers}; "
"ownership is ambiguous"
],
}
owner = candidates[0]
pr_number = owner["pr_number"]
if owner.get("head_branch") != branch_name:
return {
"status": LINKAGE_UNKNOWN,
"pr_number": pr_number,
"candidate_pr_numbers": numbers,
"reasons": [
f"PR #{pr_number} head branch '{owner.get('head_branch')}' does not "
f"match worktree branch '{branch_name}'"
],
}
if not owner.get("merged"):
return {
"status": LINKAGE_OPEN,
"pr_number": pr_number,
"candidate_pr_numbers": numbers,
"pr_head_sha": owner.get("head_sha"),
"reasons": [f"owning PR #{pr_number} is not merged"],
}
return {
"status": LINKAGE_MERGED,
"pr_number": pr_number,
"candidate_pr_numbers": numbers,
"pr_head_sha": owner.get("head_sha"),
"reasons": [],
}
def assess_merged_pr_worktree_cleanup(
*,
linkage: dict[str, Any] | None,
head_sha: str | None,
head_in_master: bool | None,
is_dirty: bool,
has_open_pr: bool,
has_active_lease: bool,
has_active_issue_lock: bool,
is_protected: bool,
has_live_session: bool = False,
) -> dict[str, Any]:
"""Decide whether a merged issue worktree satisfies the full cleanup policy.
Every condition must be independently proven: conclusive merged-PR
ownership, agreement between the worktree branch and the PR head branch,
containment of the worktree head in authoritative master (which is what
proves no unmerged commits remain), absence of any open/competing PR,
lease, issue lock, or live session, a clean tree, and a worktree that is
not the protected control checkout. Anything unknown blocks.
"""
link = linkage or {
"status": LINKAGE_UNKNOWN,
"pr_number": None,
"reasons": ["no linkage assessment supplied"],
}
status = link.get("status")
reasons: list[str] = []
if status != LINKAGE_MERGED:
reasons.extend(
link.get("reasons") or ["owning PR could not be conclusively identified"]
)
if is_protected:
reasons.append("worktree is protected or the stable control checkout")
if is_dirty:
reasons.append("worktree has uncommitted changes")
if has_open_pr:
reasons.append("worktree branch has an open PR")
if has_active_lease:
reasons.append("worktree has an active lease")
if has_active_issue_lock:
reasons.append("an active issue lock references this branch")
if has_live_session:
reasons.append("a live process or session is using this worktree")
if not head_sha:
reasons.append("worktree head sha is unknown")
if head_in_master is None:
reasons.append("containment of the worktree head in master is unknown")
elif not head_in_master:
reasons.append(
"worktree head is not contained in authoritative master "
"(unmerged commits remain)"
)
proven = not reasons
return {
"linkage_status": status,
"pr_number": link.get("pr_number"),
"pr_head_sha": link.get("pr_head_sha"),
"head_in_master": head_in_master,
"proven": proven,
"block_reasons": reasons,
}
def classify_worktree(
*,
workflow_type: str,
@@ -373,8 +181,6 @@ def classify_worktree(
ttl_expired: bool = False,
is_protected: bool = False,
metadata_known: bool = True,
merged_pr_cleanup: dict[str, Any] | None = None,
has_live_session: bool = False,
) -> str:
"""Classify a worktree, safety-first: any preservation signal wins.
@@ -393,8 +199,6 @@ def classify_worktree(
return CLASS_ACTIVE_ISSUE_WORK # never auto-deleted (criterion 8)
if has_active_issue_lock:
return CLASS_ACTIVE_ISSUE_WORK
if has_live_session:
return CLASS_ACTIVE_ISSUE_WORK # a live session still owns this tree
if not metadata_known or workflow_type == WORKFLOW_UNKNOWN:
return CLASS_UNSAFE_UNKNOWN # never auto-deleted without proof
@@ -403,15 +207,7 @@ def classify_worktree(
if is_detached or branch_gone:
return CLASS_DETACHED_REVIEW_LEFTOVER
return CLASS_CLEAN_STALE_REMOVABLE
if workflow_type == WORKFLOW_ISSUE_WORK:
# #858: an issue worktree becomes removable only on authoritative
# merged-PR evidence satisfying the whole cleanup policy. Age alone
# never proves the branch landed, so TTL cannot qualify one by itself
# — otherwise a worktree holding unmerged commits would be reclaimed.
if (merged_pr_cleanup or {}).get("proven"):
return CLASS_CLEAN_STALE_REMOVABLE
return CLASS_ACTIVE_ISSUE_WORK
# conflict_fix: only removable once the TTL has expired.
# issue_work / conflict_fix: only removable once the TTL has expired.
if ttl_expired:
return CLASS_CLEAN_STALE_REMOVABLE
return CLASS_ACTIVE_ISSUE_WORK
@@ -604,20 +400,6 @@ def remove_worktree(project_root: str, path: str) -> dict[str, Any]:
}
def head_contained_in_ref(
project_root: str, head_sha: str | None, ref: str | None
) -> bool | None:
"""Return True when ``head_sha`` is already contained in ``ref``.
Shares :mod:`merged_cleanup_reconcile`'s ancestry check so the audit and
the PR-scoped reconciler agree on what "already landed" means (#858).
Returns None when containment cannot be determined, which fails closed.
"""
if not head_sha or not ref:
return None
return is_head_ancestor_of_ref(project_root, head_sha, ref)
def _is_under_branches(project_root: str, path: str) -> bool:
branches_root = os.path.join(os.path.abspath(project_root), "branches")
return os.path.abspath(path or "").startswith(branches_root + os.sep)
@@ -631,30 +413,16 @@ def audit_branches_directory(
active_issue_branches: set[str] | None = None,
now: datetime | str | None = None,
ttl_hours: float = DEFAULT_TTL_HOURS,
pr_index: dict[str, list[dict[str, Any]]] | None = None,
leased_issue_numbers: set[int] | None = None,
live_session_paths: set[str] | None = None,
master_ref: str | None = None,
) -> dict[str, Any]:
"""Classify every session-owned worktree under ``branches/``.
Read-only: shells out to git for discovery and dirty state, then applies
the pure classifier. Returns per-worktree classifications, counts, the
list of removable candidates, and the ``git worktree list`` proof.
``pr_index`` (see :func:`build_pr_index`) supplies the authoritative PR
ownership used to link issue worktrees to their merged PR (#858).
``master_ref`` is the ref a worktree head must be contained in before it
can be considered landed. Both are optional and their absence only ever
fails closed: without them no issue worktree becomes removable.
"""
open_pr_branches = open_pr_branches or set()
leased_branches = leased_branches or set()
active_issue_branches = active_issue_branches or set()
leased_issue_numbers = leased_issue_numbers or set()
live_session_paths = {
os.path.abspath(p) for p in (live_session_paths or set()) if p
}
worktrees: list[dict[str, Any]] = []
for entry in list_worktrees(project_root):
@@ -665,42 +433,12 @@ def audit_branches_directory(
)
dirty_state = read_worktree_dirty(path)
is_dirty = bool(dirty_state.get("dirty"))
head_sha = entry.get("head")
linkage = resolve_owning_pr(branch=branch, pr_index=pr_index)
metadata = build_worktree_metadata(
path=path,
branch=branch,
head_sha=head_sha,
pr_number=linkage.get("pr_number"),
path=path, branch=branch, head_sha=entry.get("head")
)
has_open_pr = bool(branch) and branch in open_pr_branches
# A lease on issue N protects that issue's own work worktree. It must
# not incidentally protect a baseline/review scratch tree that merely
# carries the same issue marker in its name, which would change the
# classification of worktrees this policy does not own.
has_active_lease = (bool(branch) and branch in leased_branches) or (
metadata["workflow_type"] == WORKFLOW_ISSUE_WORK
and metadata.get("issue_number") is not None
and metadata["issue_number"] in leased_issue_numbers
)
has_active_lease = bool(branch) and branch in leased_branches
has_active_lock = bool(branch) and branch in active_issue_branches
has_live_session = bool(path) and os.path.abspath(path) in live_session_paths
head_in_master = (
head_contained_in_ref(project_root, head_sha, master_ref)
if master_ref
else None
)
merged_pr_cleanup = assess_merged_pr_worktree_cleanup(
linkage=linkage,
head_sha=head_sha,
head_in_master=head_in_master,
is_dirty=is_dirty,
has_open_pr=has_open_pr,
has_active_lease=has_active_lease,
has_active_issue_lock=has_active_lock,
is_protected=is_protected,
has_live_session=has_live_session,
)
ttl_expired = is_ttl_expired(
last_used_at=metadata.get("last_used_at"), now=now, ttl_hours=ttl_hours
)
@@ -714,8 +452,6 @@ def audit_branches_directory(
branch_gone=branch is None and not entry.get("detached"),
ttl_expired=ttl_expired,
is_protected=is_protected,
merged_pr_cleanup=merged_pr_cleanup,
has_live_session=has_live_session,
)
metadata["cleanup_eligibility"] = classification
worktrees.append(
@@ -727,10 +463,7 @@ def audit_branches_directory(
"has_open_pr": has_open_pr,
"has_active_lease": has_active_lease,
"has_active_issue_lock": has_active_lock,
"has_live_session": has_live_session,
"is_protected": is_protected,
"merged_pr_linkage": linkage,
"merged_pr_cleanup": merged_pr_cleanup,
"classification": classification,
"removable": is_removable(classification),
}