Compare commits

...
Author SHA1 Message Date
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
sysadmin 620ed6e9a9 Merge pull request 'fix(lock): allow exact-owner renewal of expired author issue locks (Closes #760)' (#791) from fix/issue-760-exact-owner-renewal into master 2026-07-22 01:02:39 -05:00
sysadminandClaude Opus 4.8 a30a3ce4c3 fix(lock): make the exact-owner renewal waiver survive downstream gates (#760)
Addresses review #499 on PR #791.

F1 — the renewal sanction was computed and then discarded twice.

`assess_issue_lock_worktree` waived base-equivalence only for
`recovery_sanctioned`, and `gitea_lock_issue` never passed the renewal waiver
into it. A branch being renewed always carries committed work, so it is never
base-equivalent, and #753 recovery refuses when the recorded PID is alive —
which is the defining condition of a renewal. Every real renewal was therefore
granted by the assessor and then rejected one gate later.

Thread `renewal_sanctioned` into `assess_issue_lock_worktree` alongside
`recovery_sanctioned`, waiving base-equivalence on the same grounds and nothing
else. Cleanliness is evaluated before the waiver and is never relaxed; the
assessment now reports which of the two waivers applied.

The MCP-level regression then exposed a second discard point: the duplicate-work
gate rejected the renewal with "open PR already covers issue" — the very PR the
lock being renewed already owns. Add `owning_pr_renewal_evidence`, the mirror of
`issue_lock_recovery.owning_pr_recovery_evidence` (#755), and carry it into the
gate only when renewal was granted. It re-checks that the PR, local, and remote
heads agree, so truncated or hand-built evidence cannot authorize an exemption.

Neither waiver is caller-supplied and `gitea_lock_issue` still gains no
parameter. Absolute wall-clock expiry is unchanged. No #790 heartbeat, sliding
expiry, fencing-token, or shared-lifecycle behavior is introduced, and #753
recovery behavior is untouched.

F2 — add tests/test_issue_760_mcp_renewal_path.py, 10 cases driving the native
`gitea_lock_issue` path against a real git repository and a real durable lock:
expired lease, live recorded PID, committed non-base-equivalent branch. Proves
the renewal completes, records prior and replacement lease evidence, advances
the generation exactly once, produces a live lock that satisfies
`verify_lock_for_mutation`, and does not claim dead-session recovery. Negative
companions prove a foreign claimant, foreign profile, unpublished branch, or
mismatched PR head cannot use the waiver, that a dirty worktree still fails
closed, and that base-equivalence still applies with no waiver at all.

Tests: new MCP suite 10 passed; combined #760 suites 50 passed; lock/lease
regression set 200 passed with 2 subtests; full suite 4240 passed, 11 failed,
6 skipped, 493 subtests passed — the same 11 pre-existing failures as the
master baseline at 3d0c13fa, no new failures.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01Ti8deB36iWcjHmE9cuxour
2026-07-22 01:10:48 -04:00
sysadminandClaude Opus 4.8 1a97ced133 feat(lock): allow exact-owner renewal of expired author issue locks (Closes #760)
An expired author issue lease could not be renewed by the exact session that
already owned it. `assess_same_issue_lease_conflict` computed same-owner
evidence and then returned on the expired branch before consulting it, and
`assess_expired_lock_reclaim` only permits takeover on a dead PID or a missing
worktree. Because the recorded PID is the long-lived MCP daemon rather than the
authoring task, a lease that expires under a live daemon is the ordinary case
for any author task outliving the TTL — and in that case the owner's own lock
became permanently unmodifiable through sanctioned tools.

Add `issue_lock_renewal`, a pure evidence assessor for that one case, and
evaluate its disposition before the expired foreign-takeover return.

Renewal requires an exact match of remote, org, repo, issue number, operation
type, branch, realpath-normalized worktree, claimant username, and claimant
profile; a registered worktree that exists, sits on the locked branch, and is
clean; local and remote heads that agree; and, when an owning PR exists, a PR
head that agrees too. No competing live lock, competing branch claim, or other
owning PR may exist. Any missing or contradictory evidence fails closed.

PID liveness is never authorization: it is recorded as evidence and is neither
necessary nor sufficient (AC16). Only an expired lease is ever a candidate, so
a live foreign lease stays non-recoverable (AC12) and dead-PID takeover keeps
its existing #601 conditions (AC11). Renewal is never caller-declarable — the
waiver is server-computed and `gitea_lock_issue` gains no parameter (AC14).

A sanctioned renewal records prior PID, prior expiry, replacement PID, new
expiry, claimant, and its supporting proof under `lease_renewal`, and reuses
the #772 compare-and-swap so two sessions observing the same expired lease
cannot both win.

Absolute wall-clock expiry is preserved. Sliding heartbeat renewal, fencing
tokens, and the shared cross-role lifecycle remain #790's scope and are
deliberately not implemented here; #790 stays sequenced behind this change.

Tests: 40 new cases covering the positive path, every near-match and
foreign-owner refusal, the non-candidate cases, the gate-ordering regression,
the renewal record, downstream mutation ownership, and AC14/AC17. Two #772 test
doubles now forward the new keyword.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01Ti8deB36iWcjHmE9cuxour
2026-07-22 00:39:51 -04:00
sysadmin 3d0c13fa5a Merge pull request 'fix(guard): quote-aware shell segmentation so only real daemon kills classify (Closes #787)' (#789) from fix/issue-787-kill-segment-separators into master 2026-07-21 21:40:02 -05:00
13 changed files with 3500 additions and 55 deletions
+13 -1
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:
+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`
+331 -15
View File
@@ -2007,6 +2007,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)
@@ -2017,6 +2018,7 @@ import issue_lock_provenance # noqa: E402
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 stacked_pr_support # noqa: E402
import merge_approval_gate # noqa: E402
import review_quarantine # noqa: E402 # #695 contaminated formal-review quarantine
@@ -2247,7 +2249,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,
@@ -2316,7 +2317,12 @@ def _resolve_issue_lock_for_pr(
return lock_data
def _save_issue_lock(data: dict, *, expected_generation: int | None = None) -> str:
def _save_issue_lock(
data: dict,
*,
expected_generation: int | None = None,
renewal_sanctioned: bool = False,
) -> str:
existing = issue_lock_store.load_issue_lock(
remote=str(data.get("remote") or ""),
org=str(data.get("org") or ""),
@@ -2328,7 +2334,9 @@ def _save_issue_lock(data: dict, *, expected_generation: int | None = None) -> s
raise RuntimeError(overwrite_block)
try:
return issue_lock_store.bind_session_lock(
data, expected_generation=expected_generation
data,
expected_generation=expected_generation,
renewal_sanctioned=renewal_sanctioned,
)
except Exception as e:
raise RuntimeError(f"Could not write issue lock file: {e}") from e
@@ -2448,6 +2456,79 @@ def _evaluate_issue_lock_recovery(
)
def _evaluate_issue_lock_renewal(
existing_lock: dict,
*,
issue_number: int,
branch_name: str,
worktree_path: str,
remote: str,
h: str | None,
o: str,
r: str,
git_state: dict,
) -> dict:
"""Gather evidence and decide exact-owner renewal of an expired lease (#760).
Mirrors ``_evaluate_issue_lock_recovery``: every input is durable lock state
or a live server-side observation (Gitea branch/PR inventory, git in the
declared worktree, the local lock store). Nothing is reachable from an MCP
caller's parameters, so no caller can assert its way into a renewal
(#760 AC14).
"""
renewal_auth = _auth(h)
try:
renewal_branches = api_get_all(
f"{repo_api_url(h, o, r)}/branches", renewal_auth
)
except Exception as exc:
raise RuntimeError(
f"Could not list branches to verify exact-owner lease renewal: {exc}"
)
remote_head: str | None = None
candidates: list[str] = []
for entry in renewal_branches:
entry_name = _branch_entry_name(entry)
if entry_name == branch_name:
remote_head = _branch_entry_commit_sha(entry)
if issue_lock_adoption.branch_carries_issue_marker(entry_name, issue_number):
candidates.append(entry_name)
pr_head: str | None = None
pr_number: int | None = None
for pull in _list_open_pulls(h, o, r, renewal_auth):
pull_head = pull.get("head") or {}
if str(pull_head.get("ref") or "") == branch_name:
pr_head = pull_head.get("sha")
pr_number = pull.get("number")
break
claimant = _work_lease_claimant(h)
return issue_lock_renewal.assess_exact_owner_lease_renewal(
existing_lock,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=worktree_path,
remote=remote,
org=o,
repo=r,
identity=claimant.get("username"),
profile=claimant.get("profile"),
operation_type=AUTHOR_ISSUE_WORK_LEASE,
current_branch=git_state.get("current_branch"),
porcelain_status=git_state.get("porcelain_status") or "",
worktree_exists=os.path.isdir(os.path.realpath(worktree_path)),
head_sha=git_state.get("head_sha"),
remote_head_sha=remote_head,
pr_head_sha=pr_head,
pr_number=pr_number,
competing_live_locks=issue_lock_store.list_live_locks(),
candidate_branches=candidates,
current_pid=os.getpid(),
)
def _work_lease_claimant(host: str | None) -> dict:
profile = get_profile()
username = _IDENTITY_CACHE.get(host) if host else None
@@ -2467,7 +2548,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,
@@ -2478,6 +2564,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,
}
@@ -3893,15 +3988,12 @@ def gitea_lock_issue(
existing_issue_lock = _load_existing_issue_lock(
remote=remote, org=o, repo=r, issue_number=issue_number
)
active_lease_block = issue_lock_store.assess_same_issue_lease_conflict(
existing_issue_lock,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_worktree,
operation_type=AUTHOR_ISSUE_WORK_LEASE,
)
if active_lease_block:
raise RuntimeError(active_lease_block)
# #760: the competing-lease disposition is decided below, once the worktree
# and Gitea evidence an exact-owner renewal depends on has actually been
# observed. Deciding it here — before any of that exists — is what made the
# same-owner allowance unreachable for an expired lease. The authoritative
# check still runs inside bind_session_lock under the per-issue flock, so
# moving this one later cannot widen the window for a competing writer.
# ── Stacked-PR base declaration (opt-in, #484) ──
# Normal work leaves stacked_base_branch None → master-equivalent path.
@@ -3967,6 +4059,53 @@ def gitea_lock_issue(
recovery_sanctioned = bool(
recovery_assessment and recovery_assessment.get("recovery_sanctioned")
)
# ── Exact-owner renewal of an expired lease (#760) ──
# The opposite trigger from #753 above: there the lease is unexpired and the
# PID is dead; here the lease has expired while the recording daemon — which
# is the long-lived MCP server, not the authoring task — may well still be
# up. Only an expired lease is assessed, so a live foreign lease is never a
# candidate (AC12) and dead-PID takeover keeps its existing conditions
# (AC11). A refusal never raises: it withholds the waiver and leaves the
# conflict check below to fail closed exactly as before.
renewal_assessment: dict | None = None
if (
existing_issue_lock
and existing_issue_lock.get("issue_number") == issue_number
and issue_lock_store.is_lease_expired(existing_issue_lock)
):
renewal_assessment = _evaluate_issue_lock_renewal(
existing_issue_lock,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_worktree,
remote=remote,
h=h,
o=o,
r=r,
git_state=git_state,
)
renewal_sanctioned = bool(
renewal_assessment and renewal_assessment.get("renewal_sanctioned")
)
active_lease_block = issue_lock_store.assess_same_issue_lease_conflict(
existing_issue_lock,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_worktree,
operation_type=AUTHOR_ISSUE_WORK_LEASE,
renewal_sanctioned=renewal_sanctioned,
)
if active_lease_block:
reasons = [active_lease_block]
# Name the exact missing evidence when this looked like a renewal, so a
# blocked owner sees why rather than only the generic takeover text.
if renewal_assessment and renewal_assessment.get("is_candidate"):
reasons.append(
issue_lock_renewal.format_renewal_refusal(renewal_assessment)
)
raise RuntimeError("; ".join(reasons))
# #755: a sanctioned dead-session recovery always has an owning open PR —
# that is what makes it a recovery rather than a fresh claim. Carry the
# server-derived owning-PR evidence into the duplicate-work gate below so
@@ -3978,6 +4117,14 @@ def gitea_lock_issue(
if recovery_sanctioned
else None
)
# #760: a sanctioned renewal owns its open PR for the same reason, so it
# needs the same exemption. Without it the duplicate-work gate rejects every
# renewal with "open PR already covers issue", which is the PR the lock
# being renewed already owns. Withheld unless renewal was granted.
if recovered_owning_pr is None and renewal_sanctioned:
recovered_owning_pr = issue_lock_renewal.owning_pr_renewal_evidence(
renewal_assessment
)
lock_assessment = issue_lock_worktree.assess_issue_lock_worktree(
worktree_path=resolved_worktree,
current_branch=git_state.get("current_branch"),
@@ -3986,6 +4133,10 @@ def gitea_lock_issue(
inspected_git_root=git_state.get("inspected_git_root"),
base_branch=git_state.get("base_branch"),
recovery_sanctioned=recovery_sanctioned,
# #760: without this the renewal waiver was computed and then discarded
# here — the exact-owner branch always carries commits, so it can never
# be base-equivalent, and every real renewal failed at this gate.
renewal_sanctioned=renewal_sanctioned,
)
if lock_assessment["block"]:
reasons = list(lock_assessment.get("reasons") or [])
@@ -3995,6 +4146,12 @@ def gitea_lock_issue(
reasons.append(
issue_lock_recovery.format_recovery_refusal(recovery_assessment)
)
# #760: same courtesy for a refused renewal, so an exact owner blocked
# at this gate sees which piece of ownership evidence was missing.
if renewal_assessment and renewal_assessment.get("is_candidate"):
reasons.append(
issue_lock_renewal.format_renewal_refusal(renewal_assessment)
)
raise RuntimeError(
issue_lock_worktree.format_issue_lock_worktree_error(
{**lock_assessment, "reasons": reasons}
@@ -4070,18 +4227,35 @@ def gitea_lock_issue(
recovery_assessment,
recovered_at=_work_lease_timestamp(_work_lease_now()),
)
if renewal_sanctioned and renewal_assessment:
# #760 AC9: record both sides of the transition — prior PID and expiry,
# replacement PID and new expiry — so a renewed lock is auditable and
# never reads as an original claim.
data["lease_renewal"] = issue_lock_renewal.build_renewal_record(
renewal_assessment,
renewed_at=_work_lease_timestamp(_work_lease_now()),
new_expires_at=str(work_lease.get("expires_at") or ""),
)
# #772 AC5: a recovery replaces a claim another session already owned, so
# its write is a compare-and-swap against the generation the assessment was
# made on. Two replacement sessions that both observed the same dead owner
# cannot both succeed — the second finds a moved generation and fails
# closed. Ordinary first-time claims keep the unconditional write.
# #760 uses the same compare-and-swap: a renewal also replaces a claim that
# already existed on disk, so two sessions that both observed the same
# expired lease cannot both win — the second finds a moved generation and
# fails closed.
expected_generation = (
issue_lock_store.lock_generation(existing_issue_lock)
if recovery_sanctioned
if (recovery_sanctioned or renewal_sanctioned)
else None
)
lock_file_path = _save_issue_lock(data, expected_generation=expected_generation)
lock_file_path = _save_issue_lock(
data,
expected_generation=expected_generation,
renewal_sanctioned=renewal_sanctioned,
)
lock_record = issue_lock_store.read_lock_file(lock_file_path) or data
freshness = issue_lock_store.assess_lock_freshness(lock_record)
competing = [
@@ -4150,6 +4324,16 @@ def gitea_lock_issue(
issue_number=issue_number,
branch_name=branch_name,
)
if renewal_sanctioned:
# #760 AC13: the renewal is visible in the native tool result, so an
# owner never has to inspect the lock file to confirm what happened.
result["lease_renewal"] = lock_record.get("lease_renewal")
result["message"] = (
f"Renewed the expired {AUTHOR_ISSUE_WORK_LEASE} lease on issue "
f"#{issue_number} for its exact recorded owner, on branch "
f"'{branch_name}' from worktree '{resolved_worktree}' "
"(fail-closed check complete)."
)
if agent_artifacts:
result["warnings"] = [
"Agent temp artifacts at repo root (delete before implementation): "
@@ -4158,6 +4342,138 @@ def gitea_lock_issue(
return result
@mcp.tool()
def gitea_heartbeat_issue_lock(
issue_number: int,
branch_name: str,
task_session_id: str | None = None,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
worktree_path: str | None = None,
expected_generation: int | None = None,
) -> dict:
"""Prove an owned author issue lease is still active (#790 Slice A).
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.
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.
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: 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: 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.
"""
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)
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
)
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)"
],
}
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,
)
outcome["operation"] = "legacy_rebind"
return outcome
outcome = issue_lock_store.heartbeat_session_lock(
remote=remote,
org=o,
repo=r,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_worktree,
identity=identity,
profile=profile,
task_session_id=str(task_session_id or ""),
expected_generation=expected_generation,
)
outcome["operation"] = "heartbeat"
return outcome
@mcp.tool()
def gitea_assess_work_issue_duplicate(
issue_number: int,
+481
View File
@@ -0,0 +1,481 @@
"""Exact-owner renewal of an expired author issue lease (#760).
An author issue lease carries an absolute wall-clock expiry stamped once at
lock time. The PID recorded alongside it is the long-lived MCP daemon, not the
authoring task, so a lease that expires while its daemon is still up is the
ordinary case for any author task that outlives the TTL — not an anomaly.
Before this module, that case was unreachable.
``issue_lock_store.assess_same_issue_lease_conflict`` computed same-owner
evidence and then returned on the expired branch before consulting it, and
``assess_expired_lock_reclaim`` only permits takeover on a dead PID or a
missing worktree. An exact owner whose daemon is alive and whose worktree is
present satisfied neither, so its own lock became permanently unmodifiable
through sanctioned tools.
This module is the pure evidence assessor for that one narrow case. It answers
a single question: may *this* session renew a lease it can prove it already
owns? It performs no mutation and no network I/O, and it never trusts a caller
assertion — every field is compared against durable lock state or a live
observation supplied by the caller and gathered server-side.
Deliberate boundaries:
* **Renewal is not takeover.** A refusal here never widens what
``assess_expired_lock_reclaim`` already allows; foreign expired locks keep
requiring a dead PID or missing worktree (#760 AC11), and a *live* foreign
lease stays non-recoverable by construction because only an expired lease is
ever a candidate (AC12).
* **PID liveness is never authorization.** A live recorded PID proves the
daemon is up, nothing more. It is recorded as evidence and is neither
necessary nor sufficient for renewal (AC16).
* **Absolute expiry is preserved.** Renewal issues a new absolute expiry from
the moment of the write. It does not introduce sliding heartbeat renewal,
lease generations as fencing tokens, or a shared cross-role lifecycle — that
is #790's scope and is deliberately not implemented here.
"""
from __future__ import annotations
import os
from typing import Any, Iterable, Mapping, Sequence
from issue_lock_store import AUTHOR_ISSUE_WORK_LEASE, is_lease_expired, is_process_alive
from reviewer_worktree import parse_dirty_tracked_files
# Outcome values.
RENEWAL_SANCTIONED = "RENEWAL_SANCTIONED"
NO_CANDIDATE = "NO_CANDIDATE"
REFUSED = "REFUSED"
# Durable fields a lock must carry before it can be considered at all.
REQUIRED_LOCK_FIELDS = ("issue_number", "branch_name", "worktree_path")
def _text(value: Any) -> str:
return str(value or "").strip()
def _same_realpath(left: str | None, right: str | None) -> bool:
if not left or not right:
return False
try:
return os.path.realpath(left) == os.path.realpath(right)
except OSError:
return left == right
def _lock_claimant(lock: Mapping[str, Any]) -> dict[str, Any]:
claimant = lock.get("claimant")
if not isinstance(claimant, Mapping):
lease = lock.get("work_lease")
claimant = lease.get("claimant") if isinstance(lease, Mapping) else None
return dict(claimant) if isinstance(claimant, Mapping) else {}
def _lock_lease(lock: Mapping[str, Any]) -> dict[str, Any]:
lease = lock.get("work_lease")
return dict(lease) if isinstance(lease, Mapping) else {}
def _lock_operation_type(lock: Mapping[str, Any]) -> str:
lease = _lock_lease(lock)
return _text(lease.get("operation_type")) or AUTHOR_ISSUE_WORK_LEASE
def _recorded_pid(lock: Mapping[str, Any]) -> Any:
pid = lock.get("session_pid")
if pid is None:
pid = lock.get("pid")
return pid
def _malformed_reasons(lock: Mapping[str, Any]) -> list[str]:
"""Names of durable fields that are missing or unusable."""
missing: list[str] = []
for field in REQUIRED_LOCK_FIELDS:
if not _text(lock.get(field)):
missing.append(field)
pid = _recorded_pid(lock)
if pid is None or _text(pid) == "":
missing.append("session_pid/pid")
else:
try:
if int(pid) <= 0:
missing.append("session_pid/pid")
except (TypeError, ValueError):
missing.append("session_pid/pid")
return missing
def _competing_lock_reasons(
competing_live_locks: Iterable[Mapping[str, Any]] | None,
*,
issue_number: int,
branch_name: str,
worktree_path: str,
) -> list[str]:
"""Live locks that would contend with this renewal (#760 AC7).
A live lock on the *same* issue cannot coexist with this expired lease, so
any live entry naming this issue, branch, or worktree belongs to somebody
else and refuses the renewal.
"""
reasons: list[str] = []
for entry in competing_live_locks or ():
if not isinstance(entry, Mapping):
continue
entry_issue = entry.get("issue_number")
entry_branch = _text(entry.get("branch_name"))
entry_worktree = _text(entry.get("worktree_path"))
if entry_issue == issue_number:
reasons.append(
f"a live lock already exists for issue #{issue_number} "
f"(pid {entry.get('pid')}); renewal would contend with it"
)
continue
if entry_branch and entry_branch == _text(branch_name):
reasons.append(
f"live lock for issue #{entry_issue} already holds branch "
f"'{branch_name}'"
)
if entry_worktree and _same_realpath(entry_worktree, worktree_path):
reasons.append(
f"live lock for issue #{entry_issue} already holds worktree "
f"'{worktree_path}'"
)
return reasons
def assess_exact_owner_lease_renewal(
existing_lock: Mapping[str, Any] | None,
*,
issue_number: int,
branch_name: str,
worktree_path: str,
remote: str,
org: str,
repo: str,
identity: str | None,
profile: str | None,
operation_type: str = AUTHOR_ISSUE_WORK_LEASE,
current_branch: str | None = None,
porcelain_status: str = "",
worktree_exists: bool = False,
head_sha: str | None = None,
remote_head_sha: str | None = None,
pr_head_sha: str | None = None,
pr_number: int | None = None,
competing_live_locks: Sequence[Mapping[str, Any]] | None = None,
candidate_branches: Sequence[str] | None = None,
current_pid: int | None = None,
now: Any = None,
) -> dict[str, Any]:
"""Decide whether an expired lease may be renewed by its exact owner.
Returns a disposition dict; it never raises and never mutates. A refusal
withholds permission, leaving every pre-existing guard to fail closed
exactly as before — this assessment can only ever *add* permission.
``NO_CANDIDATE`` means the situation is not an exact-owner renewal at all
(no lock, different issue, different operation, or an unexpired lease) and
the caller should carry on with its normal path. ``REFUSED`` means it looked
like one but the evidence did not hold, and ``reasons`` names exactly what
was missing.
"""
evidence: dict[str, Any] = {
"issue_number": issue_number,
"branch_name": branch_name,
"worktree_path": worktree_path,
"remote": remote,
"org": org,
"repo": repo,
"operation_type": operation_type,
"identity": identity,
"profile": profile,
}
def _result(outcome: str, reasons: list[str], **extra: Any) -> dict[str, Any]:
return {
"outcome": outcome,
"renewal_sanctioned": outcome == RENEWAL_SANCTIONED,
"is_candidate": outcome in (RENEWAL_SANCTIONED, REFUSED),
"reasons": reasons,
"evidence": {**evidence, **extra},
}
if not isinstance(existing_lock, Mapping) or not existing_lock:
return _result(NO_CANDIDATE, ["no existing lock to renew"])
if existing_lock.get("issue_number") != issue_number:
return _result(
NO_CANDIDATE,
[
f"existing lock is for issue #{existing_lock.get('issue_number')}, "
f"not #{issue_number}"
],
)
existing_operation = _lock_operation_type(existing_lock)
if existing_operation != operation_type:
return _result(
NO_CANDIDATE,
[
f"existing lease operation '{existing_operation}' is not "
f"'{operation_type}'"
],
)
# Only an *expired* lease is ever a renewal candidate. An unexpired lease —
# live, or stale by dead PID — is somebody else's problem: the first needs no
# renewal, and the second is #753's dead-session recovery. This is also what
# makes a live foreign lease non-recoverable here (#760 AC12).
if not is_lease_expired(existing_lock, now=now):
return _result(
NO_CANDIDATE,
["lease has not expired; renewal does not apply"],
)
malformed = _malformed_reasons(existing_lock)
if malformed:
return _result(
REFUSED,
["durable lock is missing or has unusable fields: " + ", ".join(malformed)],
)
lease = _lock_lease(existing_lock)
claimant = _lock_claimant(existing_lock)
recorded_pid = _recorded_pid(existing_lock)
prior_expires_at = _text(lease.get("expires_at"))
# #760 AC16: recorded purely as evidence. A live daemon PID is neither
# necessary nor sufficient for renewal, and nothing below branches on it.
recorded_pid_alive = is_process_alive(recorded_pid)
extra: dict[str, Any] = {
"prior_pid": recorded_pid,
"prior_pid_alive": recorded_pid_alive,
"prior_expires_at": prior_expires_at,
"replacement_pid": current_pid,
"recorded_claimant": claimant,
"head_sha": head_sha,
"remote_head_sha": remote_head_sha,
"pr_head_sha": pr_head_sha,
"pr_number": pr_number,
}
reasons: list[str] = []
# ── AC3: exact ownership identity ──
if _text(existing_lock.get("remote")) != _text(remote):
reasons.append(
f"recorded remote '{existing_lock.get('remote')}' does not match "
f"'{remote}'"
)
if _text(existing_lock.get("org")) != _text(org):
reasons.append(
f"recorded org '{existing_lock.get('org')}' does not match '{org}'"
)
if _text(existing_lock.get("repo")) != _text(repo):
reasons.append(
f"recorded repo '{existing_lock.get('repo')}' does not match '{repo}'"
)
if _text(existing_lock.get("branch_name")) != _text(branch_name):
reasons.append(
f"recorded branch '{existing_lock.get('branch_name')}' does not match "
f"'{branch_name}'"
)
if not _same_realpath(_text(existing_lock.get("worktree_path")), worktree_path):
reasons.append(
f"recorded worktree '{existing_lock.get('worktree_path')}' does not "
f"match '{worktree_path}'"
)
recorded_identity = _text(claimant.get("username"))
recorded_profile = _text(claimant.get("profile"))
if not recorded_identity or not recorded_profile:
reasons.append(
"durable lock does not record both a claimant username and profile"
)
if recorded_identity and recorded_identity != _text(identity):
reasons.append(
f"recorded claimant '{recorded_identity}' does not match active "
f"identity '{_text(identity) or 'unknown'}'"
)
if recorded_profile and recorded_profile != _text(profile):
reasons.append(
f"recorded profile '{recorded_profile}' does not match active profile "
f"'{_text(profile) or 'unknown'}'"
)
# ── AC4: the registered worktree still exists, is on the branch, and is clean ──
if not worktree_exists:
reasons.append(f"declared worktree '{worktree_path}' does not exist")
if _text(current_branch) != _text(branch_name):
reasons.append(
f"worktree is on branch '{_text(current_branch) or 'unknown'}', not "
f"'{branch_name}'"
)
dirty = parse_dirty_tracked_files(porcelain_status or "")
if dirty:
reasons.append(
"worktree has uncommitted tracked changes: " + ", ".join(sorted(dirty))
)
# ── AC5/AC6: published heads must agree ──
if not _text(head_sha):
reasons.append("local head could not be observed")
if not _text(remote_head_sha):
reasons.append(
"remote branch head could not be observed; an unpublished branch "
"cannot prove exact-owner renewal"
)
if _text(head_sha) and _text(remote_head_sha) and head_sha != remote_head_sha:
reasons.append(
f"local head {head_sha} does not equal remote head {remote_head_sha}"
)
if pr_number is not None:
if not _text(pr_head_sha):
reasons.append(f"owning PR #{pr_number} head could not be observed")
elif _text(head_sha) and pr_head_sha != head_sha:
reasons.append(
f"owning PR #{pr_number} head {pr_head_sha} does not equal local "
f"head {head_sha}"
)
# ── AC7: nothing else claims this work ──
reasons.extend(
_competing_lock_reasons(
competing_live_locks,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=worktree_path,
)
)
other_branches = [
name
for name in (candidate_branches or ())
if _text(name) and _text(name) != _text(branch_name)
]
if other_branches:
reasons.append(
"other branches already carry this issue marker: "
+ ", ".join(sorted(other_branches))
)
if reasons:
return _result(REFUSED, reasons, **extra)
return _result(
RENEWAL_SANCTIONED,
[
f"exact owner '{recorded_identity}' ({recorded_profile}) proved "
f"ownership of issue #{issue_number} on branch '{branch_name}' from "
f"worktree '{worktree_path}'; local, remote"
+ (f", and PR #{pr_number}" if pr_number is not None else "")
+ f" heads all equal {head_sha}; lease expired at "
f"{prior_expires_at or 'unknown'}"
],
**extra,
)
def owning_pr_renewal_evidence(
assessment: Mapping[str, Any] | None,
) -> dict[str, Any] | None:
"""Server-derived proof of the open PR a sanctioned renewal already owns.
The mirror of ``issue_lock_recovery.owning_pr_recovery_evidence`` (#755) for
the renewal disposition. An exact-owner renewal of a published branch is, by
construction, renewal of work that already has an open PR — so the
duplicate-work gate's linked-open-PR blocker would otherwise discard every
sanctioned renewal, exactly as it once discarded every sanctioned recovery.
Returns ``None`` unless renewal was actually granted and the evidence names
one owning PR whose head agrees with both the local and remote heads the
assessor accepted. Nothing is caller-supplied: every field is copied from
evidence built out of durable lock state plus live git/Gitea observation.
Renewal has no descendant case — it requires the local, remote, and PR heads
to be equal — so there is only one head to report.
"""
if not isinstance(assessment, Mapping):
return None
if assessment.get("outcome") != RENEWAL_SANCTIONED:
return None
if not assessment.get("renewal_sanctioned"):
return None
evidence = assessment.get("evidence") or {}
branch_name = _text(evidence.get("branch_name"))
pr_head = _text(evidence.get("pr_head_sha"))
local_head = _text(evidence.get("head_sha"))
remote_head = _text(evidence.get("remote_head_sha"))
raw_pr_number = evidence.get("pr_number")
if raw_pr_number is None or not branch_name or not pr_head:
return None
# The assessor already required these to agree. Re-check, so a truncated or
# hand-built evidence map can never authorize an exemption.
if pr_head != local_head or pr_head != remote_head:
return None
try:
pr_number = int(raw_pr_number)
issue_number = int(evidence.get("issue_number"))
except (TypeError, ValueError):
return None
return {
"issue_number": issue_number,
"pr_number": pr_number,
"branch_name": branch_name,
"head_sha": pr_head,
"recorded_head": pr_head,
"accepted_head": pr_head,
"head_relation": "equal",
}
def build_renewal_record(
assessment: Mapping[str, Any] | None,
*,
renewed_at: str,
new_expires_at: str,
) -> dict[str, Any]:
"""Durable audit record for a sanctioned renewal (#760 AC9).
Records both sides of the transition — prior PID and expiry, replacement PID
and new expiry — so a renewed lock is never mistakable for an original
claim, and so the evidence the waiver was granted on stays inspectable.
"""
data = dict(assessment or {})
evidence = dict(data.get("evidence") or {})
recorded_claimant = dict(evidence.get("recorded_claimant") or {})
return {
"renewed": bool(data.get("renewal_sanctioned")),
"renewed_at": renewed_at,
"prior_pid": evidence.get("prior_pid"),
"prior_pid_alive": evidence.get("prior_pid_alive"),
"prior_expires_at": evidence.get("prior_expires_at"),
"replacement_pid": evidence.get("replacement_pid"),
"new_expires_at": new_expires_at,
"identity": recorded_claimant.get("username"),
"profile": recorded_claimant.get("profile"),
"branch_name": evidence.get("branch_name"),
"worktree_path": evidence.get("worktree_path"),
"head_sha": evidence.get("head_sha"),
"remote_head_sha": evidence.get("remote_head_sha"),
"pr_head_sha": evidence.get("pr_head_sha"),
"pr_number": evidence.get("pr_number"),
"reason": "expired lease renewed by its exact recorded owner",
"proof": list(data.get("reasons") or []),
}
def format_renewal_refusal(assessment: Mapping[str, Any] | None) -> str:
"""One-line refusal summary for a blocked caller."""
data = dict(assessment or {})
reasons = list(data.get("reasons") or [])
if not reasons:
return "exact-owner lease renewal was not available (no evidence recorded)"
return "exact-owner lease renewal refused: " + "; ".join(reasons)
+576 -30
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._+-]+")
@@ -168,6 +180,7 @@ def bind_session_lock(
lock_dir: str | None = None,
*,
expected_generation: int | None = None,
renewal_sanctioned: bool = False,
) -> str:
"""Persist a keyed lock and bind it to the current process session.
@@ -220,6 +233,7 @@ def bind_session_lock(
issue_number=issue_number,
branch_name=str(record.get("branch_name") or ""),
worktree_path=str(record.get("worktree_path") or ""),
renewal_sanctioned=renewal_sanctioned,
)
if lease_block:
raise RuntimeError(lease_block)
@@ -251,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))
@@ -334,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
@@ -354,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:
@@ -444,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")
@@ -483,9 +1009,19 @@ def assess_same_issue_lease_conflict(
branch_name: str,
worktree_path: str,
operation_type: str = AUTHOR_ISSUE_WORK_LEASE,
renewal_sanctioned: bool = False,
now: datetime | None = None,
) -> str | None:
"""Return a fail-closed error when a competing live lease blocks acquisition."""
"""Return a fail-closed error when a competing live lease blocks acquisition.
``renewal_sanctioned`` is set only when
``issue_lock_renewal.assess_exact_owner_lease_renewal`` has already proven,
from the durable lock plus live server-side observation, that this session
is the exact recorded owner of an *expired* lease (#760). It is never a
caller-supplied parameter of any MCP tool (#760 AC14): the server computes
it and passes it down. Left False, every pre-existing disposition is
unchanged.
"""
if not existing_lock:
return None
@@ -506,6 +1042,16 @@ def assess_same_issue_lease_conflict(
and _same_realpath(str(existing_worktree or ""), worktree_path)
)
if is_lease_expired(existing_lock, now=now):
# #760 AC1/AC2: exact-owner renewal is a different disposition from
# foreign takeover and is evaluated first. Before this, both branches
# below returned unconditionally, so the same_owner allowance further
# down was unreachable for every expired lease — an owner could never
# renew its own lock once the wall clock passed, no matter how complete
# its ownership evidence. Requires BOTH the locally recomputed
# same_owner match and the server-proven renewal waiver; either alone is
# insufficient.
if same_owner and renewal_sanctioned:
return None
reclaim = assess_expired_lock_reclaim(existing_lock, now=now)
if reclaim.get("reclaim_allowed"):
# #601: expired + dead pid / missing worktree may be reclaimed
+25 -3
View File
@@ -285,6 +285,7 @@ def assess_issue_lock_worktree(
base_branch: str | None = None,
base_branches: frozenset[str] | None = None,
recovery_sanctioned: bool = False,
renewal_sanctioned: bool = False,
) -> dict:
"""Fail closed when lock preconditions are not met on the declared worktree.
@@ -296,6 +297,19 @@ def assess_issue_lock_worktree(
by construction and could never satisfy it. Every other precondition —
notably worktree cleanliness — still applies unchanged, and brand-new issue
claims keep the full base-equivalence requirement.
``renewal_sanctioned`` waives base-equivalence on exactly the same grounds
for the other proven-ownership case (#760): ``issue_lock_renewal`` has shown
that an *expired* lease is being renewed by its exact recorded owner — same
remote, org, repo, issue, operation, branch, realpath-normalized worktree,
claimant username and profile — with the local head matching the remote head
and any owning PR head. Such a branch carries committed work for the same
reason a recovered one does, so it can never be base-equivalent either.
Both waivers relax this one requirement and nothing else. Neither is
caller-supplied: each is computed server-side from durable lock state plus
live observation. With both False every precondition applies exactly as
before.
"""
bases = base_branches or BASE_BRANCHES
reasons: list[str] = []
@@ -314,9 +328,12 @@ def assess_issue_lock_worktree(
f"(dirty files: {', '.join(dirty_files)})"
)
if recovery_sanctioned:
if recovery_sanctioned or renewal_sanctioned:
# Base-equivalence intentionally not evaluated: ownership was proven
# against the durable lock record instead (#753).
# against the durable lock record instead — by dead-session recovery
# (#753) or by exact-owner renewal of an expired lease (#760). Every
# other precondition above and below still applies; cleanliness in
# particular is checked before this branch and is never waived.
pass
elif base_equivalent is False:
reasons.append(
@@ -347,6 +364,7 @@ def assess_issue_lock_worktree(
base_branch=base_branch,
base_equivalent=base_equivalent,
recovery_sanctioned=recovery_sanctioned,
renewal_sanctioned=renewal_sanctioned,
)
@@ -406,6 +424,7 @@ def _assessment(
base_branch: str | None = None,
base_equivalent: bool | None = None,
recovery_sanctioned: bool = False,
renewal_sanctioned: bool = False,
) -> dict:
return {
"proven": proven,
@@ -418,7 +437,10 @@ def _assessment(
"base_branch": base_branch,
"base_equivalent": base_equivalent,
"recovery_sanctioned": recovery_sanctioned,
"base_equivalence_waived": bool(recovery_sanctioned),
"renewal_sanctioned": renewal_sanctioned,
# Either proven-ownership waiver relaxes base-equivalence; the two are
# reported separately so an audit can tell which one applied.
"base_equivalence_waived": bool(recovery_sanctioned or renewal_sanctioned),
}
+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,
}
+9
View File
@@ -32,6 +32,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.issue.comment",
"role": "author",
},
# #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",
},
"set_issue_labels": {
"permission": "gitea.issue.comment",
"role": "author",
@@ -0,0 +1,447 @@
"""Exact-owner renewal of an expired author issue lease (#760).
Covers the renewal disposition that lets the exact recorded owner re-acquire
its own lock after the wall-clock lease expires — including while the recording
MCP daemon PID is still alive — plus every rejection condition that must keep
failing closed, and the pre-existing dead-PID and live-foreign dispositions
that must remain untouched.
"""
import inspect
import os
import subprocess
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
sys.path.insert(0, str(__import__("pathlib").Path(__file__).resolve().parent.parent))
import issue_lock_renewal # noqa: E402
import issue_lock_store # noqa: E402
ISSUE = 5150
BRANCH = f"fix/issue-{ISSUE}-demo"
WORKTREE = "/scratch/wt-5150"
HEAD = "c" * 40
OTHER_SHA = "d" * 40
IDENTITY = "example-user"
PROFILE = "example-author"
REMOTE = "prgs"
ORG = "ExampleOrg"
REPO = "ExampleRepo"
def dead_pid() -> int:
"""A PID that has certainly exited (spawned, then reaped)."""
proc = subprocess.Popen([sys.executable, "-c", "pass"])
proc.wait()
return proc.pid
def past_ts(hours: int = 1) -> str:
return (
(datetime.now(timezone.utc) - timedelta(hours=hours))
.isoformat()
.replace("+00:00", "Z")
)
def future_ts(hours: int = 4) -> str:
return (
(datetime.now(timezone.utc) + timedelta(hours=hours))
.isoformat()
.replace("+00:00", "Z")
)
def make_lock(*, expires_at: str | None = None, pid: int | None = None, **overrides):
"""An expired lock owned by a still-alive daemon PID — the #760 condition."""
lock = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": WORKTREE,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
# os.getpid() is unambiguously alive: the whole point of #760 is that
# daemon liveness is not evidence of an active author task.
"session_pid": os.getpid() if pid is None else pid,
"lock_generation": 3,
"work_lease": {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"branch": BRANCH,
"worktree_path": WORKTREE,
"claimant": {"username": IDENTITY, "profile": PROFILE},
"created_at": past_ts(5),
"expires_at": expires_at or past_ts(),
},
}
lease_overrides = overrides.pop("work_lease", None)
if lease_overrides:
lock["work_lease"].update(lease_overrides)
lock.update(overrides)
return lock
def assess(lock=None, **overrides):
"""Run the assessor with all-passing evidence unless overridden."""
kwargs = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": WORKTREE,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"identity": IDENTITY,
"profile": PROFILE,
"current_branch": BRANCH,
"porcelain_status": "",
"worktree_exists": True,
"head_sha": HEAD,
"remote_head_sha": HEAD,
"pr_head_sha": None,
"pr_number": None,
"competing_live_locks": [],
"candidate_branches": [BRANCH],
"current_pid": 4242,
}
kwargs.update(overrides)
return issue_lock_renewal.assess_exact_owner_lease_renewal(
make_lock() if lock is None else lock, **kwargs
)
class ExactOwnerRenewalGranted(unittest.TestCase):
"""AC1/AC3-AC7: the positive path."""
def test_expired_lease_alive_pid_exact_owner_is_renewable(self):
result = assess()
self.assertEqual(result["outcome"], issue_lock_renewal.RENEWAL_SANCTIONED)
self.assertTrue(result["renewal_sanctioned"])
self.assertTrue(result["is_candidate"])
def test_renewal_holds_when_owning_pr_head_matches(self):
result = assess(pr_number=999, pr_head_sha=HEAD)
self.assertTrue(result["renewal_sanctioned"])
def test_evidence_records_both_sides_of_the_transition(self):
result = assess()
evidence = result["evidence"]
self.assertEqual(evidence["prior_pid"], os.getpid())
self.assertTrue(evidence["prior_pid_alive"])
self.assertEqual(evidence["replacement_pid"], 4242)
self.assertTrue(evidence["prior_expires_at"])
class ExactOwnerRenewalRefused(unittest.TestCase):
"""AC3-AC8: every near-match must fail closed, one reason at a time."""
def _refused(self, **overrides):
result = assess(**overrides)
self.assertEqual(result["outcome"], issue_lock_renewal.REFUSED)
self.assertFalse(result["renewal_sanctioned"])
self.assertTrue(result["reasons"])
return result
def test_different_branch_refused(self):
result = self._refused(branch_name=f"fix/issue-{ISSUE}-other")
self.assertTrue(any("branch" in r for r in result["reasons"]))
def test_different_worktree_refused(self):
result = self._refused(worktree_path="/scratch/somewhere-else")
self.assertTrue(any("worktree" in r for r in result["reasons"]))
def test_different_claimant_refused(self):
result = self._refused(identity="someone-else")
self.assertTrue(any("claimant" in r for r in result["reasons"]))
def test_different_profile_refused(self):
result = self._refused(profile="other-author")
self.assertTrue(any("profile" in r for r in result["reasons"]))
def test_different_remote_org_or_repo_refused(self):
self._refused(remote="dadeschools")
self._refused(org="OtherOrg")
self._refused(repo="OtherRepo")
def test_dirty_worktree_refused(self):
result = self._refused(porcelain_status=" M gitea_mcp_server.py\n")
self.assertTrue(any("uncommitted" in r for r in result["reasons"]))
def test_missing_worktree_refused(self):
result = self._refused(worktree_exists=False)
self.assertTrue(any("does not exist" in r for r in result["reasons"]))
def test_worktree_on_wrong_branch_refused(self):
self._refused(current_branch="master")
def test_local_and_remote_head_mismatch_refused(self):
result = self._refused(remote_head_sha=OTHER_SHA)
self.assertTrue(
any("does not equal remote head" in r for r in result["reasons"])
)
def test_unpublished_branch_refused(self):
result = self._refused(remote_head_sha=None)
self.assertTrue(any("remote branch head" in r for r in result["reasons"]))
def test_pr_head_mismatch_refused(self):
result = self._refused(pr_number=999, pr_head_sha=OTHER_SHA)
self.assertTrue(any("does not equal local" in r for r in result["reasons"]))
def test_unobservable_pr_head_refused(self):
self._refused(pr_number=999, pr_head_sha=None)
def test_competing_live_lock_on_same_issue_refused(self):
result = self._refused(
competing_live_locks=[
{"issue_number": ISSUE, "branch_name": BRANCH, "pid": 777}
]
)
self.assertTrue(any("live lock" in r for r in result["reasons"]))
def test_competing_live_lock_holding_the_branch_refused(self):
self._refused(
competing_live_locks=[
{"issue_number": 111, "branch_name": BRANCH, "worktree_path": ""}
]
)
def test_competing_branch_claim_refused(self):
result = self._refused(candidate_branches=[BRANCH, f"feat/issue-{ISSUE}-rival"])
self.assertTrue(any("issue marker" in r for r in result["reasons"]))
def test_malformed_durable_lock_refused(self):
lock = make_lock()
lock["worktree_path"] = ""
result = assess(lock)
self.assertEqual(result["outcome"], issue_lock_renewal.REFUSED)
def test_lock_without_recorded_claimant_refused(self):
lock = make_lock()
lock["work_lease"]["claimant"] = {}
result = assess(lock)
self.assertEqual(result["outcome"], issue_lock_renewal.REFUSED)
class NotARenewalCandidate(unittest.TestCase):
"""AC12 and scope: situations renewal must decline to judge at all."""
def test_live_foreign_lease_is_never_a_candidate(self):
lock = make_lock(expires_at=future_ts())
result = assess(lock, identity="someone-else")
self.assertEqual(result["outcome"], issue_lock_renewal.NO_CANDIDATE)
self.assertFalse(result["renewal_sanctioned"])
def test_unexpired_lease_is_never_a_candidate(self):
lock = make_lock(expires_at=future_ts())
result = assess(lock)
self.assertEqual(result["outcome"], issue_lock_renewal.NO_CANDIDATE)
def test_dead_pid_under_unexpired_lease_stays_with_753(self):
"""The opposite trigger; #760 must not re-own it."""
lock = make_lock(expires_at=future_ts(), pid=dead_pid())
result = assess(lock)
self.assertEqual(result["outcome"], issue_lock_renewal.NO_CANDIDATE)
def test_absent_lock_is_not_a_candidate(self):
result = assess({})
self.assertEqual(result["outcome"], issue_lock_renewal.NO_CANDIDATE)
def test_different_issue_is_not_a_candidate(self):
lock = make_lock()
lock["issue_number"] = ISSUE + 1
result = assess(lock)
self.assertEqual(result["outcome"], issue_lock_renewal.NO_CANDIDATE)
def test_different_operation_type_is_not_a_candidate(self):
lock = make_lock()
lock["work_lease"]["operation_type"] = "review_pr_work"
result = assess(lock)
self.assertEqual(result["outcome"], issue_lock_renewal.NO_CANDIDATE)
class DaemonPidIsNotTaskLiveness(unittest.TestCase):
"""AC16: a live recorded PID is never, by itself, authorization."""
def test_alive_pid_alone_does_not_authorize_renewal(self):
# Every ownership fact except the live PID is wrong.
result = assess(identity="someone-else", branch_name="fix/issue-1-nope")
self.assertEqual(result["outcome"], issue_lock_renewal.REFUSED)
self.assertTrue(result["evidence"]["prior_pid_alive"])
def test_renewal_does_not_require_a_dead_pid(self):
result = assess()
self.assertTrue(result["evidence"]["prior_pid_alive"])
self.assertTrue(result["renewal_sanctioned"])
def test_dead_pid_does_not_block_an_otherwise_exact_owner(self):
lock = make_lock(pid=dead_pid())
result = assess(lock)
self.assertTrue(result["renewal_sanctioned"])
class ConflictGateOrdering(unittest.TestCase):
"""AC2: the same-owner allowance is reachable on an expired lease.
These cases need a worktree that genuinely exists on disk. The #601 reclaim
affordance already permits takeover when the recorded worktree is missing,
so a fictional path would satisfy the gate for the wrong reason and never
exercise the ordering defect this issue is about.
"""
@classmethod
def setUpClass(cls):
cls._tmp = tempfile.TemporaryDirectory()
cls.worktree = cls._tmp.name
@classmethod
def tearDownClass(cls):
cls._tmp.cleanup()
def present_lock(self, **overrides):
return make_lock(worktree_path=self.worktree, **overrides)
def test_expired_same_owner_is_allowed_when_renewal_is_sanctioned(self):
block = issue_lock_store.assess_same_issue_lease_conflict(
self.present_lock(),
issue_number=ISSUE,
branch_name=BRANCH,
worktree_path=self.worktree,
renewal_sanctioned=True,
)
self.assertIsNone(block)
def test_expired_same_owner_still_blocks_without_the_waiver(self):
"""Regression for the ordering defect: no waiver, no change in behavior.
Live PID and a present worktree, so the #601 reclaim affordance refuses;
before #760 this was the permanent dead end for an exact owner.
"""
lock = self.present_lock()
self.assertFalse(
issue_lock_store.assess_expired_lock_reclaim(lock)["reclaim_allowed"]
)
block = issue_lock_store.assess_same_issue_lease_conflict(
lock,
issue_number=ISSUE,
branch_name=BRANCH,
worktree_path=self.worktree,
)
self.assertIsNotNone(block)
self.assertIn("Recovery review is required", block)
def test_waiver_does_not_unlock_a_different_owner(self):
"""AC11: the waiver is scoped by same_owner, not merely by its own flag."""
block = issue_lock_store.assess_same_issue_lease_conflict(
self.present_lock(),
issue_number=ISSUE,
branch_name=f"fix/issue-{ISSUE}-someone-else",
worktree_path=self.worktree,
renewal_sanctioned=True,
)
self.assertIsNotNone(block)
self.assertIn("Recovery review is required", block)
def test_live_lease_disposition_is_unchanged(self):
"""AC12: a live foreign lease still blocks, waiver or not."""
block = issue_lock_store.assess_same_issue_lease_conflict(
self.present_lock(expires_at=future_ts()),
issue_number=ISSUE,
branch_name=f"fix/issue-{ISSUE}-someone-else",
worktree_path="/scratch/other",
renewal_sanctioned=True,
)
self.assertIsNotNone(block)
self.assertIn("already has an active", block)
def test_dead_pid_reclaim_path_is_unchanged(self):
"""AC11: expired + dead PID still reclaims through the #601 affordance."""
lock = self.present_lock(pid=dead_pid())
reclaim = issue_lock_store.assess_expired_lock_reclaim(lock)
self.assertTrue(reclaim["reclaim_allowed"])
block = issue_lock_store.assess_same_issue_lease_conflict(
lock,
issue_number=ISSUE,
branch_name=BRANCH,
worktree_path=self.worktree,
)
self.assertIsNone(block)
class RenewalRecordAndDownstream(unittest.TestCase):
"""AC9/AC10: durable audit trail, and a renewed lock that actually works."""
def test_record_captures_prior_and_replacement_state(self):
assessment = assess()
record = issue_lock_renewal.build_renewal_record(
assessment,
renewed_at="2026-01-01T00:00:00Z",
new_expires_at="2026-01-01T04:00:00Z",
)
self.assertTrue(record["renewed"])
self.assertEqual(record["prior_pid"], os.getpid())
self.assertEqual(record["new_expires_at"], "2026-01-01T04:00:00Z")
self.assertEqual(record["renewed_at"], "2026-01-01T00:00:00Z")
self.assertEqual(record["identity"], IDENTITY)
self.assertEqual(record["profile"], PROFILE)
self.assertTrue(record["prior_expires_at"])
self.assertTrue(record["proof"])
def test_renewed_lock_satisfies_verify_lock_for_mutation(self):
renewed = make_lock(expires_at=future_ts())
renewed["session_pid"] = os.getpid()
renewed["lease_renewal"] = {"renewed": True}
verdict = issue_lock_store.verify_lock_for_mutation(
renewed,
issue_number=ISSUE,
branch_name=BRANCH,
)
self.assertTrue(verdict["proven"])
self.assertFalse(verdict["block"])
def test_refusal_message_names_the_missing_evidence(self):
assessment = assess(porcelain_status=" M gitea_mcp_server.py\n")
message = issue_lock_renewal.format_renewal_refusal(assessment)
self.assertIn("refused", message)
self.assertIn("uncommitted", message)
class NoCallerControlledRenewalFlag(unittest.TestCase):
"""AC14: renewal eligibility is never declarable by a caller."""
def test_lock_issue_tool_exposes_no_renewal_parameter(self):
import gitea_mcp_server
target = gitea_mcp_server.gitea_lock_issue
target = getattr(target, "fn", getattr(target, "__wrapped__", target))
params = set(inspect.signature(target).parameters)
for forbidden in ("renewal_sanctioned", "renew", "allow_renewal", "is_owner"):
self.assertNotIn(forbidden, params)
def test_store_defaults_to_no_waiver(self):
params = inspect.signature(
issue_lock_store.assess_same_issue_lease_conflict
).parameters
self.assertIs(params["renewal_sanctioned"].default, False)
bind_params = inspect.signature(issue_lock_store.bind_session_lock).parameters
self.assertIs(bind_params["renewal_sanctioned"].default, False)
class NoIssueNumberSpecialCasing(unittest.TestCase):
"""AC17: no repository issue or PR number is special-cased."""
def test_module_contains_no_hardcoded_issue_special_cases(self):
source = inspect.getsource(issue_lock_renewal)
code = "\n".join(
line for line in source.splitlines() if not line.strip().startswith("#")
)
for literal in ("757", "759", "760"):
self.assertNotIn(f"== {literal}", code)
self.assertNotIn(f"issue_number == {literal}", code)
if __name__ == "__main__":
unittest.main()
+342
View File
@@ -0,0 +1,342 @@
"""MCP-level exact-owner lease renewal through ``gitea_lock_issue`` (#760).
The unit suite in ``test_issue_760_exact_owner_lease_renewal`` proves the
renewal *disposition*. It cannot prove the disposition survives the rest of the
tool, and it did not: the waiver was computed and then discarded before
``assess_issue_lock_worktree``, so every real renewal still failed on
base-equivalence. A branch being renewed always carries committed work, so it is
never base-equivalent by construction — exactly the argument #753 already makes
for recovery.
These tests drive the public tool end to end against a real git repository and a
real durable lock file, composing every gate in the production order.
"""
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 mcp_server # noqa: E402
ISSUE = 9760
BRANCH = f"fix/issue-{ISSUE}-renewal-mcp"
IDENTITY = "example-user"
PROFILE = "test-author-prgs"
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
def _past_ts(hours: int = 1) -> str:
return (
(datetime.now(timezone.utc) - timedelta(hours=hours))
.isoformat()
.replace("+00:00", "Z")
)
class _RenewalMcpBase(unittest.TestCase):
"""Real git repo + durable expired lock owned by a live PID.
The recorded PID is ``os.getpid()`` — unambiguously alive. That is the whole
point of #760: the PID belongs to the long-lived MCP daemon, so its liveness
says nothing about whether the authoring task still holds the work.
"""
def setUp(self):
self.lock_dir = tempfile.TemporaryDirectory()
self.addCleanup(self.lock_dir.cleanup)
self.repo = tempfile.mkdtemp(prefix="issue760-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()
# The branch carries committed work, so it is NOT base-equivalent.
self._git("checkout", "-q", "-b", BRANCH)
with open(os.path.join(self.repo, "work.txt"), "w") as fh:
fh.write("author work\n")
self._git("add", "work.txt")
self._git("commit", "-q", "-m", "author work")
self.head_sha = self._git("rev-parse", "HEAD").stdout.strip()
self.worktree = os.path.realpath(self.repo)
def write_expired_lock(self, **overrides):
path = issue_lock_store.lock_file_path(
remote="prgs",
org=ORG,
repo=REPO,
issue_number=ISSUE,
lock_dir=self.lock_dir.name,
)
claimant = {"username": IDENTITY, "profile": PROFILE}
pid = overrides.pop("session_pid", os.getpid())
overrides.pop("pid", None)
lease_overrides = overrides.pop("work_lease", {})
data = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"remote": "prgs",
"org": ORG,
"repo": REPO,
"worktree_path": self.worktree,
"session_pid": pid,
"pid": pid,
"lock_generation": 3,
"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": _past_ts(5),
"last_heartbeat_at": _past_ts(5),
"expires_at": _past_ts(), # already expired
},
"lock_provenance": issue_lock_provenance.build_sanctioned_lock_provenance(
tool="gitea_lock_issue",
claimant=claimant,
),
}
data["work_lease"].update(lease_overrides)
data.update(overrides)
data["session_pid"] = pid
data["pid"] = pid
data["lock_file_path"] = path
issue_lock_store.save_lock_file(path, data)
return path
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="", branch=BRANCH, head=None):
return {
"current_branch": branch,
"porcelain_status": porcelain,
# The decisive fact: a branch carrying work is never base-equivalent.
"base_equivalent": False,
"head_sha": head or 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,
):
"""Drive the public tool for the published exact-owner renewal shape."""
if branch_entries is None:
branch_entries = [{"name": BRANCH, "commit": {"id": self.head_sha}}]
if open_prs is None:
open_prs = [{"number": 4242, "head": {"ref": BRANCH, "sha": self.head_sha}}]
if git_state is None:
git_state = 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,
)
class TestRenewalReachableThroughTool(_RenewalMcpBase):
"""F1: the sanctioned renewal must survive every downstream gate."""
def test_expired_lease_live_pid_exact_owner_renews_through_the_tool(self):
prior = issue_lock_store.read_lock_file(self.write_expired_lock())
self.assertTrue(issue_lock_store.is_lease_expired(prior))
self.assertTrue(issue_lock_store.is_process_alive(prior["session_pid"]))
result = self.run_lock_issue()
self.assertTrue(result["success"], result)
self.assertEqual(result["issue_number"], ISSUE)
self.assertEqual(result["branch_name"], BRANCH)
# The renewal is reported natively, so no lock-file inspection is needed.
self.assertIn("lease_renewal", result)
self.assertTrue(result["lease_renewal"]["renewed"])
self.assertIn("Renewed the expired", result["message"])
def test_renewed_lock_records_prior_and_replacement_evidence(self):
prior = issue_lock_store.read_lock_file(self.write_expired_lock())
prior_expiry = prior["work_lease"]["expires_at"]
prior_generation = issue_lock_store.lock_generation(prior)
result = self.run_lock_issue()
written = issue_lock_store.read_lock_file(result["lock_file_path"])
renewal = written["lease_renewal"]
self.assertTrue(renewal["renewed"])
self.assertEqual(renewal["prior_pid"], prior["session_pid"])
self.assertTrue(renewal["prior_pid_alive"])
self.assertEqual(renewal["prior_expires_at"], prior_expiry)
self.assertEqual(renewal["identity"], IDENTITY)
self.assertEqual(renewal["profile"], PROFILE)
self.assertEqual(renewal["head_sha"], self.head_sha)
self.assertTrue(renewal["proof"])
# New expiry is a fresh absolute stamp, later than the one it replaced.
self.assertEqual(renewal["new_expires_at"], written["work_lease"]["expires_at"])
self.assertGreater(renewal["new_expires_at"], prior_expiry)
# Compare-and-swap advanced the generation exactly once.
self.assertEqual(
issue_lock_store.lock_generation(written), prior_generation + 1
)
def test_renewed_lock_is_live_and_satisfies_mutation_ownership(self):
self.write_expired_lock()
result = self.run_lock_issue()
written = issue_lock_store.read_lock_file(result["lock_file_path"])
self.assertTrue(issue_lock_store.assess_lock_freshness(written)["live"])
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 test_recovery_record_is_not_written_for_a_live_owner_renewal(self):
"""#753 recovery must not be claimed when the recorded PID is alive."""
self.write_expired_lock()
result = self.run_lock_issue()
written = issue_lock_store.read_lock_file(result["lock_file_path"])
self.assertNotIn("dead_session_recovery", written)
class TestRenewalWaiverIsNarrow(_RenewalMcpBase):
"""The waiver relaxes base-equivalence and nothing else."""
def test_dirty_worktree_still_blocks_a_would_be_renewal(self):
"""Cleanliness is never waived; the renewal assessor refuses first.
A dirty worktree makes the renewal refuse, so no waiver is issued and
the lease-conflict gate fails closed ahead of the worktree gate. The
refusal names the uncommitted files, so the owner still learns why.
"""
self.write_expired_lock()
with self.assertRaises(Exception) as ctx:
self.run_lock_issue(
git_state=self._git_state(porcelain=" M gitea_mcp_server.py\n")
)
message = str(ctx.exception)
self.assertIn("Recovery review is required before takeover", message)
self.assertIn("worktree has uncommitted tracked changes", message)
self.assertIn("gitea_mcp_server.py", message)
def test_foreign_claimant_cannot_use_the_waiver(self):
"""A near-match owner gets no renewal and no base-equivalence waiver."""
self.write_expired_lock()
with self.assertRaises(Exception) as ctx:
self.run_lock_issue(identity="someone-else")
message = str(ctx.exception)
self.assertIn("Recovery review is required before takeover", message)
# The refusal names the missing ownership evidence (#760 diagnostics).
self.assertIn("does not match active identity", message)
def test_foreign_profile_cannot_use_the_waiver(self):
self.write_expired_lock()
with self.assertRaises(Exception) as ctx:
self.run_lock_issue(profile="other-author")
self.assertIn(
"Recovery review is required before takeover", str(ctx.exception)
)
def test_unpublished_branch_cannot_use_the_waiver(self):
"""No remote head to agree with, so exact-owner renewal is refused."""
self.write_expired_lock()
with self.assertRaises(Exception) as ctx:
self.run_lock_issue(branch_entries=[], open_prs=[])
self.assertIn(
"Recovery review is required before takeover", str(ctx.exception)
)
def test_pr_head_mismatch_cannot_use_the_waiver(self):
self.write_expired_lock()
other = "9" * 40
with self.assertRaises(Exception) as ctx:
self.run_lock_issue(
open_prs=[{"number": 4242, "head": {"ref": BRANCH, "sha": other}}]
)
self.assertIn(
"Recovery review is required before takeover", str(ctx.exception)
)
def test_non_base_equivalent_branch_still_blocks_without_any_waiver(self):
"""No durable lock at all: the ordinary base-equivalence rule applies."""
with self.assertRaises(Exception) as ctx:
self.run_lock_issue()
self.assertIn("must be base-equivalent", str(ctx.exception))
if __name__ == "__main__":
unittest.main()
@@ -871,9 +871,13 @@ class TestAc6McpUnpublishedClaimRecovery(_UnpublishedMcpBase):
save_calls: list[dict] = []
real_save = mcp_server._save_issue_lock
def tracking_save(data, *, expected_generation=None):
def tracking_save(data, *, expected_generation=None, renewal_sanctioned=False):
save_calls.append({"expected_generation": expected_generation, "data": dict(data)})
return real_save(data, expected_generation=expected_generation)
return real_save(
data,
expected_generation=expected_generation,
renewal_sanctioned=renewal_sanctioned,
)
with patch("mcp_server._save_issue_lock", side_effect=tracking_save):
result = self.run_lock_issue()
@@ -926,18 +930,33 @@ class TestAc6McpUnpublishedClaimRecovery(_UnpublishedMcpBase):
real_bind = issue_lock_store.bind_session_lock
bind_calls: list[int | None] = []
def racing_bind(data, lock_dir=None, expected_generation=None):
# #760 added the renewal waiver keyword; the double forwards it verbatim
# so this race still exercises the real compare-and-swap.
def racing_bind(
data, lock_dir=None, expected_generation=None, renewal_sanctioned=False
):
bind_calls.append(expected_generation)
if expected_generation is None:
return real_bind(data, lock_dir=lock_dir, expected_generation=None)
return real_bind(
data,
lock_dir=lock_dir,
expected_generation=None,
renewal_sanctioned=renewal_sanctioned,
)
# First concurrent writer wins.
if len([c for c in bind_calls if c is not None]) == 1:
return real_bind(
data, lock_dir=lock_dir, expected_generation=expected_generation
data,
lock_dir=lock_dir,
expected_generation=expected_generation,
renewal_sanctioned=renewal_sanctioned,
)
# Second concurrent writer still holds the pre-race generation.
return real_bind(
data, lock_dir=lock_dir, expected_generation=expected_generation
data,
lock_dir=lock_dir,
expected_generation=expected_generation,
renewal_sanctioned=renewal_sanctioned,
)
# First recovery succeeds and advances generation.
+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()