diff --git a/branch_cleanup_guard.py b/branch_cleanup_guard.py index 64db84a..084a4f1 100644 --- a/branch_cleanup_guard.py +++ b/branch_cleanup_guard.py @@ -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: diff --git a/docs/mcp-tool-inventory.md b/docs/mcp-tool-inventory.md index a44a6c7..351b054 100644 --- a/docs/mcp-tool-inventory.md +++ b/docs/mcp-tool-inventory.md @@ -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` diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index 0ac7b4b..39d4094 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -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) @@ -2248,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, @@ -2548,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, @@ -2559,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, } @@ -4328,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, diff --git a/issue_lock_store.py b/issue_lock_store.py index 713fa2a..698cca2 100644 --- a/issue_lock_store.py +++ b/issue_lock_store.py @@ -15,15 +15,27 @@ import json import os import re import tempfile +import uuid from contextlib import contextmanager from datetime import datetime, timedelta, timezone from typing import Any +import lease_policy + LOCK_DIR_ENV = "GITEA_ISSUE_LOCK_DIR" DEFAULT_LOCK_DIR = os.path.expanduser("~/.cache/gitea-tools/issue-locks") -WORK_LEASE_TTL_HOURS = 4 AUTHOR_ISSUE_WORK_LEASE = "author_issue_work" +# Freshness classifications. ``STATUS_STALE`` remains the dead-PID band that +# #753 recovery keys on; the two bands below are new in #790 Slice A and apply +# only to leases minted under the heartbeat lifecycle. +STATUS_LIVE = "live" +STATUS_EXPIRED = "expired" +STATUS_ABSENT = "absent" +STATUS_STALE = "stale" +STATUS_STALE_MISSED_HEARTBEAT = "stale_missed_heartbeat" +STATUS_STALE_ABSOLUTE_CAP = "stale_absolute_cap" + _SAFE_SEGMENT_RE = re.compile(r"[^A-Za-z0-9._+-]+") @@ -253,6 +265,331 @@ def bind_session_lock( return path +def _ownership_refusals( + lock: dict[str, Any], + *, + issue_number: int, + branch_name: str, + worktree_path: str, + identity: str | None, + profile: str | None, +) -> list[str]: + """Exact-ownership mismatches between a durable lock and a live caller. + + Shared by the heartbeat writer and the legacy rebind path so the two cannot + disagree about what "the same owner" means. Every field is compared against + durable state; nothing is taken on the caller's word beyond the identity the + server itself resolved. + """ + reasons: list[str] = [] + if lock.get("issue_number") != issue_number: + reasons.append( + f"lock targets issue #{lock.get('issue_number')}, not #{issue_number}" + ) + if str(lock.get("branch_name") or "") != str(branch_name or ""): + reasons.append( + f"lock branch '{lock.get('branch_name')}' does not match '{branch_name}'" + ) + if not _same_realpath(str(lock.get("worktree_path") or ""), worktree_path): + reasons.append( + f"lock worktree '{lock.get('worktree_path')}' does not match " + f"'{worktree_path}'" + ) + lease = lock.get("work_lease") if isinstance(lock, dict) else None + claimant = lease.get("claimant") if isinstance(lease, dict) else None + claimant = claimant if isinstance(claimant, dict) else {} + recorded_identity = str(claimant.get("username") or "").strip() + recorded_profile = str(claimant.get("profile") or "").strip() + if not recorded_identity or not recorded_profile: + reasons.append("lock does not record both a claimant username and profile") + if recorded_identity and recorded_identity != str(identity or "").strip(): + reasons.append( + f"lock claimant '{recorded_identity}' does not match active identity " + f"'{str(identity or '').strip() or 'unknown'}'" + ) + if recorded_profile and recorded_profile != str(profile or "").strip(): + reasons.append( + f"lock profile '{recorded_profile}' does not match active profile " + f"'{str(profile or '').strip() or 'unknown'}'" + ) + return reasons + + +def _refusal(reasons: list[str], **extra: Any) -> dict[str, Any]: + return {"success": False, "performed": False, "reasons": reasons, **extra} + + +def heartbeat_session_lock( + *, + remote: str, + org: str, + repo: str, + issue_number: int, + branch_name: str, + worktree_path: str, + identity: str | None, + profile: str | None, + task_session_id: str, + expected_generation: int | None = None, + lock_dir: str | None = None, + now: datetime | None = None, +) -> dict[str, Any]: + """Slide a heartbeat-lifecycle lease forward (#790 Slice A, A4). + + The write happens inside the same per-issue ``flock`` that serializes + acquisition, and under the #772 generation compare-and-swap, so a heartbeat + can never race a concurrent reclaim: whichever lands first moves the + generation and the other fails closed. + + Refuses — never revives — in every ambiguous case. A lease that has already + lapsed past its grace is *not* heartbeatable: allowing that would let a + session that stopped proving liveness restore ownership retroactively, which + is precisely the revival AC-N5 forbids. Such a session must go through the + sanctioned reclaim path, which mints a fresh generation. + """ + current = _lease_now(now) + root = _ensure_lock_dir(lock_dir) + path = lock_file_path( + remote=remote, org=org, repo=repo, issue_number=issue_number, lock_dir=root + ) + declared_session = str(task_session_id or "").strip() + if not declared_session: + return _refusal(["no task_session_id supplied (fail closed)"]) + + sentinel = flock_path(path) + try: + with _exclusive_file_lock(sentinel): + lock = read_lock_file(path) + if not lock: + return _refusal([f"no durable lock for issue #{issue_number}"]) + + if is_legacy_lease(lock): + return _refusal( + [ + "lock predates the heartbeat lifecycle; it must be rebound " + "by its exact owner before it can be heartbeated" + ], + lifecycle=lease_lifecycle_version(lock), + legacy_lease=True, + ) + + reasons = _ownership_refusals( + lock, + issue_number=issue_number, + branch_name=branch_name, + worktree_path=worktree_path, + identity=identity, + profile=profile, + ) + recorded_session = lease_task_session_id(lock) + if not recorded_session: + reasons.append( + "lock declares the heartbeat lifecycle but records no " + "task_session_id (fail closed)" + ) + elif recorded_session != declared_session: + # A superseded session holding an old identifier cannot heartbeat + # over the session that replaced it. + reasons.append( + "task_session_id does not match the session recorded on the lock" + ) + if reasons: + return _refusal(reasons) + + current_generation = lock_generation(lock) + if ( + expected_generation is not None + and current_generation != expected_generation + ): + return _refusal( + [ + f"lock generation changed: expected {expected_generation}, " + f"found {current_generation}; another session reclaimed or " + "replaced this claim (fail closed)" + ], + lock_generation=current_generation, + ) + + freshness = assess_lock_freshness(lock, now=current) + if not freshness.get("live"): + return _refusal( + [ + f"lease is not live ({freshness.get('status')}): " + f"{freshness.get('reason')}; a lapsed lease must be " + "reclaimed, not heartbeated" + ], + freshness=freshness, + ) + + policy = lease_policy.policy_for(lease_task_class(lock)) + expires = current + timedelta(minutes=policy.initial_ttl_minutes) + record = dict(lock) + lease = dict(record.get("work_lease") or {}) + prior_heartbeat = lease.get("last_heartbeat_at") + lease["last_heartbeat_at"] = _format_lease_timestamp(current) + lease["expires_at"] = _format_lease_timestamp(expires) + try: + lease["heartbeat_count"] = int(lease.get("heartbeat_count") or 0) + 1 + except (TypeError, ValueError): + lease["heartbeat_count"] = 1 + record["work_lease"] = lease + record["lock_generation"] = current_generation + 1 + save_lock_file(path, record) + except LockContentionError as exc: + return _refusal([f"issue #{issue_number} lock contention: {exc} (fail closed)"]) + + return { + "success": True, + "performed": True, + "issue_number": issue_number, + "branch_name": branch_name, + "worktree_path": worktree_path, + "task_session_id": declared_session, + "lock_generation": record["lock_generation"], + "prior_generation": current_generation, + "prior_heartbeat_at": prior_heartbeat, + "last_heartbeat_at": lease["last_heartbeat_at"], + "expires_at": lease["expires_at"], + "heartbeat_count": lease["heartbeat_count"], + "lock_file_path": path, + "policy": lease_policy.describe(lease_task_class(record)), + "freshness": assess_lock_freshness(record, now=current), + } + + +def rebind_legacy_lock( + *, + remote: str, + org: str, + repo: str, + issue_number: int, + branch_name: str, + worktree_path: str, + identity: str | None, + profile: str | None, + expected_generation: int | None = None, + lock_dir: str | None = None, + now: datetime | None = None, +) -> dict[str, Any]: + """Move a legacy lock into the heartbeat lifecycle (#790 AC-N8). + + One of the two sanctioned exits from the preserved-expiry legacy state; the + other is terminal retirement, which is Slice B. Only the exact recorded + owner may rebind, and only while the legacy lock is still live under its + original absolute expiry — an already-expired legacy lease belongs to the + #760 renewal path or #601 reclaim, and this must not become a second, weaker + way to revive one. + + The rebind mints a genuine task-session identifier and a genuine first + heartbeat. It does not fabricate history: the original creation and expiry + are preserved under ``legacy_origin`` for audit, and the new lifecycle's + absolute cap runs from the rebind, not from the legacy claim. + """ + current = _lease_now(now) + root = _ensure_lock_dir(lock_dir) + path = lock_file_path( + remote=remote, org=org, repo=repo, issue_number=issue_number, lock_dir=root + ) + sentinel = flock_path(path) + try: + with _exclusive_file_lock(sentinel): + lock = read_lock_file(path) + if not lock: + return _refusal([f"no durable lock for issue #{issue_number}"]) + if not is_legacy_lease(lock): + return _refusal( + [ + "lock is already on the heartbeat lifecycle; use the " + "heartbeat path" + ], + lifecycle=lease_lifecycle_version(lock), + legacy_lease=False, + ) + + reasons = _ownership_refusals( + lock, + issue_number=issue_number, + branch_name=branch_name, + worktree_path=worktree_path, + identity=identity, + profile=profile, + ) + if reasons: + return _refusal(reasons) + + current_generation = lock_generation(lock) + if ( + expected_generation is not None + and current_generation != expected_generation + ): + return _refusal( + [ + f"lock generation changed: expected {expected_generation}, " + f"found {current_generation} (fail closed)" + ], + lock_generation=current_generation, + ) + + freshness = assess_lock_freshness(lock, now=current) + if not freshness.get("live"): + return _refusal( + [ + f"legacy lease is not live ({freshness.get('status')}): " + f"{freshness.get('reason')}; rebinding is not a recovery " + "path for a lapsed lease" + ], + freshness=freshness, + ) + + policy = lease_policy.policy_for(lease_task_class(lock)) + expires = current + timedelta(minutes=policy.initial_ttl_minutes) + session_id = mint_task_session_id(lease_task_class(lock)) + record = dict(lock) + lease = dict(record.get("work_lease") or {}) + legacy_origin = { + "created_at": lease.get("created_at"), + "expires_at": lease.get("expires_at"), + "last_heartbeat_at": lease.get("last_heartbeat_at"), + "lifecycle": lease_policy.LIFECYCLE_LEGACY, + } + lease["lifecycle_version"] = lease_policy.LIFECYCLE_HEARTBEAT_V1 + lease["task_session_id"] = session_id + lease["created_at"] = _format_lease_timestamp(current) + lease["last_heartbeat_at"] = _format_lease_timestamp(current) + lease["expires_at"] = _format_lease_timestamp(expires) + lease["heartbeat_count"] = 1 + record["work_lease"] = lease + record["legacy_rebind"] = { + "rebound_at": _format_lease_timestamp(current), + "task_session_id": session_id, + "prior_generation": current_generation, + "legacy_origin": legacy_origin, + "reason": ( + "legacy lock rebound into the heartbeat lifecycle by its exact " + "recorded owner" + ), + } + record["lock_generation"] = current_generation + 1 + save_lock_file(path, record) + except LockContentionError as exc: + return _refusal([f"issue #{issue_number} lock contention: {exc} (fail closed)"]) + + return { + "success": True, + "performed": True, + "issue_number": issue_number, + "task_session_id": session_id, + "lock_generation": record["lock_generation"], + "prior_generation": current_generation, + "lifecycle": lease_policy.LIFECYCLE_HEARTBEAT_V1, + "legacy_rebind": record["legacy_rebind"], + "expires_at": lease["expires_at"], + "last_heartbeat_at": lease["last_heartbeat_at"], + "lock_file_path": path, + "freshness": assess_lock_freshness(record, now=current), + } + + def read_session_issue_lock(lock_dir: str | None = None) -> dict[str, Any] | None: root = (lock_dir or default_lock_dir()).strip() pointer = read_lock_file(session_pointer_path(root)) @@ -336,6 +673,16 @@ def _parse_lease_timestamp(value: str | None) -> datetime | None: return None +def _format_lease_timestamp(value: datetime) -> str: + """Serialize a lease timestamp in the durable ``...Z`` form already on disk.""" + return ( + value.astimezone(timezone.utc) + .replace(microsecond=0) + .isoformat() + .replace("+00:00", "Z") + ) + + def lease_expires_at(lock: dict[str, Any] | None) -> datetime | None: if not lock: return None @@ -356,60 +703,216 @@ def is_lease_live(lock: dict[str, Any] | None, *, now: datetime | None = None) - return assess_lock_freshness(lock, now=now)["live"] +def lease_task_class(lock_data: dict[str, Any] | None) -> str: + """Policy task class for a durable lock; author work when unrecorded.""" + lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None + if isinstance(lease, dict): + recorded = str(lease.get("operation_type") or "").strip() + if recorded: + return recorded + return AUTHOR_ISSUE_WORK_LEASE + + +def lease_lifecycle_version(lock_data: dict[str, Any] | None) -> str: + """Read the durable lifecycle marker (#790 AC-N8). + + The marker is the *only* discriminator between a heartbeat-lifecycle lease + and a legacy one. Timestamps are deliberately not consulted: a lock minted + before this lifecycle existed has ``last_heartbeat_at == created_at`` + forever, and reading that equality as "recently heartbeated" would treat + every never-heartbeated legacy lock as fresh — the precise inversion AC-N8 + forbids. A newly minted heartbeat lease also has the two equal, so the + equality carries no information in either direction. + """ + lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None + if isinstance(lease, dict): + recorded = str(lease.get("lifecycle_version") or "").strip() + if recorded: + return recorded + return lease_policy.LIFECYCLE_LEGACY + + +def is_legacy_lease(lock_data: dict[str, Any] | None) -> bool: + """True when a lock predates the shared heartbeat lifecycle.""" + return lease_lifecycle_version(lock_data) != lease_policy.LIFECYCLE_HEARTBEAT_V1 + + +def lease_task_session_id(lock_data: dict[str, Any] | None) -> str: + """Recorded per-task session identifier, or empty for a legacy lock.""" + lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None + if isinstance(lease, dict): + return str(lease.get("task_session_id") or "").strip() + return "" + + +def mint_task_session_id(task_class: str = AUTHOR_ISSUE_WORK_LEASE) -> str: + """Mint an ownership key for one task (#790 AC-N1). + + Deliberately contains no process identifier. The recorded PID belongs to the + long-lived MCP daemon, which outlives any individual task and is reused by + every task it serves, so PID digits cannot identify *which* task holds a + claim. The PID is still recorded alongside this value as evidence. + """ + prefix = _sanitize_segment(str(task_class or AUTHOR_ISSUE_WORK_LEASE)) + return f"{prefix}-{uuid.uuid4().hex[:16]}" + + +def _lease_heartbeat_at(lock_data: dict[str, Any] | None) -> datetime | None: + lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None + heartbeat_at = None + if isinstance(lock_data, dict): + heartbeat_at = _parse_lease_timestamp(lock_data.get("last_heartbeat_at")) + if heartbeat_at is None and isinstance(lease, dict): + heartbeat_at = _parse_lease_timestamp(lease.get("last_heartbeat_at")) + return heartbeat_at + + def assess_lock_freshness( lock_data: dict[str, Any] | None, *, now: datetime | None = None, ) -> dict[str, Any]: - """Classify a lock as live, expired, stale, or absent.""" + """Classify a lock as live, expired, stale, or absent. + + #790 Slice A makes the heartbeat load-bearing. Before this change + ``last_heartbeat_at`` was parsed and then never consulted: liveness was + decided entirely by the absolute ``expires_at`` and by PID liveness, and + since the recorded PID is the long-lived MCP daemon, an abandoned author + task stayed "live" for the full four-hour TTL. + + Two rules govern the rewrite: + + * **An alive PID never establishes freshness** (AC-N2). It proves the daemon + is up, nothing about the task. It is recorded as evidence and no branch + returns ``live`` because of it. + * **A dead PID still corroborates staleness.** The dead-PID band is + unchanged and still precedes every heartbeat evaluation, so #753 + dead-session recovery keys on exactly the classification it always did. + + Legacy leases (AC-N8) keep their recorded absolute expiry and are never + evaluated against the short heartbeat grace, so deploying this change cannot + make an existing claim instantly reclaimable. + """ current = _lease_now(now) if not lock_data: return { - "status": "absent", + "status": STATUS_ABSENT, "live": False, "stale": False, "reason": "no lock record", } - expires_at = lease_expires_at(lock_data) lease = lock_data.get("work_lease") - heartbeat_at = _parse_lease_timestamp(lock_data.get("last_heartbeat_at")) - if heartbeat_at is None and isinstance(lease, dict): - heartbeat_at = _parse_lease_timestamp(lease.get("last_heartbeat_at")) + expires_at = lease_expires_at(lock_data) + heartbeat_at = _lease_heartbeat_at(lock_data) + created_at = ( + _parse_lease_timestamp(lease.get("created_at")) + if isinstance(lease, dict) + else None + ) pid = lock_data.get("session_pid") if pid is None: pid = lock_data.get("pid") + # Evidence only. Never consulted to grant liveness (AC-N2). pid_alive = is_process_alive(pid) if pid is not None else False - if expires_at and expires_at <= current: - return { - "status": "expired", - "live": False, - "stale": True, - "reason": f"lease expired at {expires_at.isoformat()}", - "pid_alive": pid_alive, - } + lifecycle = lease_lifecycle_version(lock_data) + legacy = lifecycle != lease_policy.LIFECYCLE_HEARTBEAT_V1 + policy = lease_policy.policy_for(lease_task_class(lock_data)) - if pid is not None and not pid_alive: - return { - "status": "stale", - "live": False, - "stale": True, - "reason": f"owner pid {pid} is not alive", - "pid_alive": False, - } - - return { - "status": "live", - "live": True, - "stale": False, - "reason": "lock heartbeat and lease are fresh", + evidence: dict[str, Any] = { "pid_alive": pid_alive, + "lifecycle": lifecycle, + "legacy_lease": legacy, + "task_session_id": lease_task_session_id(lock_data) or None, "heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None, "expires_at": expires_at.isoformat() if expires_at else None, } + def _result(status: str, *, live: bool, reason: str, **extra: Any) -> dict[str, Any]: + return { + "status": status, + "live": live, + "stale": not live and status != STATUS_ABSENT, + "reason": reason, + **evidence, + **extra, + } + + if legacy: + # AC-N8: the preserved absolute expiry is the only clock for a lock + # written before task-session heartbeats existed. + if expires_at and expires_at <= current: + return _result( + STATUS_EXPIRED, + live=False, + reason=f"lease expired at {expires_at.isoformat()}", + ) + if pid is not None and not pid_alive: + return _result( + STATUS_STALE, live=False, reason=f"owner pid {pid} is not alive" + ) + return _result( + STATUS_LIVE, + live=True, + reason=( + "legacy lease is within its recorded absolute expiry; the " + "heartbeat grace does not apply retroactively" + ), + legacy_expiry_preserved=True, + ) + + # ── Heartbeat lifecycle ── + if pid is not None and not pid_alive: + # Unchanged dead-PID band: #753 recovery depends on this exact status. + return _result(STATUS_STALE, live=False, reason=f"owner pid {pid} is not alive") + + if heartbeat_at is None: + # Contradictory: a heartbeat lease must carry a heartbeat. Fail closed. + return _result( + STATUS_STALE_MISSED_HEARTBEAT, + live=False, + reason=( + f"lease declares lifecycle '{lifecycle}' but records no " + "last_heartbeat_at (fail closed)" + ), + ) + + if policy.absolute_cap_hours and created_at is not None: + cap_at = created_at + timedelta(hours=policy.absolute_cap_hours) + if cap_at <= current: + return _result( + STATUS_STALE_ABSOLUTE_CAP, + live=False, + reason=( + f"lease exceeded its {policy.absolute_cap_hours}h absolute cap " + f"at {cap_at.isoformat()}; canonical re-adoption is required" + ), + absolute_cap_at=cap_at.isoformat(), + ) + + grace_at = heartbeat_at + timedelta(minutes=policy.missed_heartbeat_grace_minutes) + if grace_at <= current or (expires_at is not None and expires_at <= current): + return _result( + STATUS_STALE_MISSED_HEARTBEAT, + live=False, + reason=( + f"no valid heartbeat since {heartbeat_at.isoformat()}; the " + f"{policy.missed_heartbeat_grace_minutes}min grace lapsed at " + f"{grace_at.isoformat()}" + ), + missed_heartbeat_since=grace_at.isoformat(), + ) + + warning_at = heartbeat_at + timedelta(minutes=policy.stale_warning_minutes) + return _result( + STATUS_LIVE, + live=True, + reason="lease heartbeat is fresh within the configured grace", + heartbeat_warning=warning_at <= current, + ) + def _same_realpath(left: str | None, right: str | None) -> bool: if not left or not right: @@ -446,6 +949,27 @@ def assess_expired_lock_reclaim( "reasons": ["lock is still live; cannot reclaim (fail closed)"], "freshness": freshness, } + status = str(freshness.get("status") or "") + if status in (STATUS_STALE_MISSED_HEARTBEAT, STATUS_STALE_ABSOLUTE_CAP): + # #790: under the heartbeat lifecycle the heartbeat *is* the liveness + # proof, so a session that stopped heartbeating past its grace has + # released its claim by definition. Requiring a dead PID on top of that + # would reinstate the original defect — the recorded PID is the shared + # daemon, which stays alive across every abandoned task it ever served. + # + # This band is unreachable for a legacy lease (AC-N8), so no lock + # written before this lifecycle can be reclaimed by this path. + return { + "reclaim_allowed": True, + "reasons": [ + f"heartbeat-lifecycle lease is {status}: {freshness.get('reason')}" + ], + "freshness": freshness, + "prior_branch": existing_lock.get("branch_name"), + "prior_worktree": existing_lock.get("worktree_path"), + "prior_pid": existing_lock.get("session_pid") or existing_lock.get("pid"), + "prior_task_session_id": lease_task_session_id(existing_lock) or None, + } pid = existing_lock.get("session_pid") if pid is None: pid = existing_lock.get("pid") diff --git a/lease_policy.py b/lease_policy.py new file mode 100644 index 0000000..e778cd3 --- /dev/null +++ b/lease_policy.py @@ -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, + } diff --git a/task_capability_map.py b/task_capability_map.py index 23409e0..fba85da 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -32,6 +32,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "permission": "gitea.issue.comment", "role": "author", }, + # #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", diff --git a/tests/test_issue_790_heartbeat_mcp_path.py b/tests/test_issue_790_heartbeat_mcp_path.py new file mode 100644 index 0000000..57e0d5b --- /dev/null +++ b/tests/test_issue_790_heartbeat_mcp_path.py @@ -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", "test@example.com") + 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() diff --git a/tests/test_issue_790_lease_policy.py b/tests/test_issue_790_lease_policy.py new file mode 100644 index 0000000..34d52d6 --- /dev/null +++ b/tests/test_issue_790_lease_policy.py @@ -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()