Compare commits

..
Author SHA1 Message Date
jcwalker3 6868b345ee Merge branch 'master' into feat/issue-634-readonly-system-health-api 2026-07-23 01:12:52 -05:00
jcwalker3 da6a864463 Merge branch 'master' into feat/issue-634-readonly-system-health-api 2026-07-23 00:06:06 -05:00
jcwalker3andClaude Opus 4.8 5494696227 feat(webui): read-only system-health API (Closes #634)
Adds `GET /api/v1/system/health`, a structured read-only health surface for
automated readiness checks, and keeps `/health` as the cheap liveness probe.

webui/system_health.py composes a DTO from fail-soft dependency probes: the
control-plane database, the local checkout, and — opt-in via `?deep=1` — live
Gitea reachability, each carrying status, reason, and probe latency. Required
probes drive readiness; the optional Gitea probe can only degrade overall
status, because local inventory stays serveable when the remote is
unreachable. A probe that did not run leaves readiness incomplete rather than
silently passing.

Read-only throughout: the control-plane database is opened through a `mode=ro`
URI because `ControlPlaneDB.__init__` creates directories and runs migrations,
which a health check must never do. No restart or reload control is exposed;
those are Phase 2 and #630 forbids process-kill recovery.

No unproven claims: `stale_runtime.mutation_safe` is true only when the
runtime, checkout, and remote commits are all known and equal, and MCP
namespaces always report `unproven` because a web process cannot exercise the
IDE-managed client path (#543). Probe details are redacted at the browser
boundary — URLs lose userinfo and query strings, credential-shaped text is
masked.

`/health` is expanded additively: every MVP key is retained, plus `started_at`,
`uptime_seconds`, and a pointer to the versioned API. The versioned route
returns 503 when not ready so automation can branch on the status code alone.

Verified at master 9eb0f29: focused file 40 passed / 11 subtests; `-k "webui or
health"` 230 passed / 159 subtests; full suite 4358 passed with the 11
pre-existing master-drift failures unchanged from the clean-master baseline.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-22 15:59:30 -05:00
12 changed files with 1328 additions and 1975 deletions
+1 -13
View File
@@ -163,19 +163,7 @@ _TERMINAL_OWNERSHIP_STATUSES = frozenset(
{"released", "abandoned", "done", "blocked", "terminal", "closed"} {"released", "abandoned", "done", "blocked", "terminal", "closed"}
) )
_EXPIRED_STATUSES = frozenset({"expired"}) _EXPIRED_STATUSES = frozenset({"expired"})
_STALE_STATUSES = frozenset( _STALE_STATUSES = frozenset({"stale", "stale_dead_process", "stale_missing_worktree"})
{
"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: def _norm_str(value: Any) -> str:
-1
View File
@@ -100,7 +100,6 @@ that gates each call, not which tools exist.
- `gitea_get_profile` - `gitea_get_profile`
- `gitea_get_runtime_context` - `gitea_get_runtime_context`
- `gitea_get_shell_health` - `gitea_get_shell_health`
- `gitea_heartbeat_issue_lock`
- `gitea_heartbeat_reviewer_pr_lease` - `gitea_heartbeat_reviewer_pr_lease`
- `gitea_inspect_workflow_lease` - `gitea_inspect_workflow_lease`
- `gitea_issue_irrecoverable_provenance_authorization` - `gitea_issue_irrecoverable_provenance_authorization`
+81 -1
View File
@@ -52,7 +52,8 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| Path | Description | | Path | Description |
|------|-------------| |------|-------------|
| `/` | Home / operator overview | | `/` | Home / operator overview |
| `/health` | JSON liveness (`status`, `service`, `mode`, `timestamp`) | | `/health` | JSON liveness (`status`, `service`, `mode`, `timestamp`, `uptime_seconds`) |
| `/api/v1/system/health` | Structured read-only system health (#634) |
| `/queue` | Live PR and issue queue dashboard (#429) | | `/queue` | Live PR and issue queue dashboard (#429) |
| `/api/queue` | JSON queue export with pagination metadata | | `/api/queue` | JSON queue export with pagination metadata |
| `/projects` | Project registry list with status and onboarding progress (#427, #635) | | `/projects` | Project registry list with status and onboarding progress (#427, #635) |
@@ -78,6 +79,85 @@ Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for `read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
local validator preview only (no Gitea mutations, no server-side storage). local validator preview only (no Gitea mutations, no server-side storage).
## System health API (#634)
`GET /api/v1/system/health` is the structured, read-only health surface for
automated readiness checks. It is the first console API under the `/api/v1`
prefix; the unversioned MVP exports remain as compatibility aliases.
`/health` is unchanged for existing consumers — every MVP key is still present
— and now also carries `started_at`, `uptime_seconds`, and a
`system_health_api` pointer. It stays deliberately cheap and runs no dependency
probe, because answering readiness costs real work.
**Status codes.** `200` when ready, `503` when a required dependency failed or
was never probed. Automation can branch on the code without parsing the body.
**Query flags.** The Gitea check is a network call, so it is opt-in:
`GET /api/v1/system/health?deep=1` runs it and caches the result for
`WEBUI_HEALTH_PROBE_TTL_SECONDS` (default 15s) so dashboard polling does not
amplify into remote load. Without the flag that probe reports `skipped`.
**Dependencies.** `control_plane_db` and `repository` are required and drive
readiness. `gitea` is optional: when it fails the overall `status` degrades but
`readiness.ready` stays true, because local inventory is still serveable. Each
entry carries `status`, `detail`, `required`, and `latency_ms`.
Two honesty rules are worth knowing before reading the payload:
* `stale_runtime.mutation_safe` is true only when the runtime, checkout, and
remote-tracking commits are all known and equal. An unfetched remote is
reported as indeterminate, never as safe.
* `mcp_namespaces` entries are always `unproven`. A web process runs outside
the IDE-managed MCP client and cannot invoke a namespace tool, so per #543
only a `client_namespace` probe can prove that path.
Sample response (abridged, healthy):
```json
{
"status": "ok",
"service": "mcp-control-plane-webui",
"mode": "read-only",
"api": "/api/v1/system/health",
"timestamp": "2026-07-22T11:04:18.512034+00:00",
"readiness": { "ready": true, "complete": true, "reasons": [] },
"version": {
"git_sha": "620ed6e9a9550b8da2ceb82d9ab8744e8920490f",
"git_describe": "v1.1.0-898-g620ed6e",
"control_plane_schema_version": 4,
"python_version": "3.14.5",
"known": true
},
"process": { "started_at": "2026-07-22T10:58:02.114+00:00", "uptime_seconds": 376.4 },
"deep_probes_requested": false,
"dependencies": [
{
"name": "control_plane_db",
"kind": "sqlite",
"status": "ok",
"detail": "schema v4 readable",
"required": true,
"healthy": true,
"latency_ms": 1.482,
"metadata": { "schema_version": 4, "active_leases": 3 }
},
{ "name": "repository", "kind": "git", "status": "ok", "required": true, "healthy": true },
{ "name": "gitea", "kind": "http", "status": "skipped", "required": false, "healthy": false }
],
"mcp_namespaces": [
{ "namespace": "gitea-author", "required_tool": "gitea_whoami", "status": "unproven" }
],
"stale_runtime": { "stale": false, "determinable": true, "mutation_safe": true, "reasons": [] },
"probe_errors": []
}
```
No restart, reload, or process-kill control is exposed here: those are Phase 2
at the earliest, and #630 forbids process-kill recovery outright. Every probe
opens its subject read-only — the control-plane database is opened through a
`mode=ro` URI so a health check can never create or migrate a schema.
## Report audit (#431) ## Report audit (#431)
Paste an LLM final report at `/audit` or POST JSON to `/api/audit`. The UI Paste an LLM final report at `/audit` or POST JSON to `/api/audit`. The UI
+2 -148
View File
@@ -2020,7 +2020,6 @@ import allocator_dependencies # noqa: E402
import dependency_graph # noqa: E402 # #784 durable dependency edges import dependency_graph # noqa: E402 # #784 durable dependency edges
import control_plane_db # noqa: E402 import control_plane_db # noqa: E402
import lease_lifecycle # noqa: E402 import lease_lifecycle # noqa: E402
import lease_policy # noqa: E402
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
import incident_bridge # noqa: E402 import incident_bridge # noqa: E402
import sentry_observability # noqa: E402 (#606 optional Sentry observability) import sentry_observability # noqa: E402 (#606 optional Sentry observability)
@@ -2263,6 +2262,7 @@ import canonical_comment_validator as ccv # noqa: E402
# GITEA_ISSUE_LOCK_DIR, bound to the current MCP session via a per-PID pointer. # 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. # Legacy global path retained only for test/doc references — do not seed manually.
ISSUE_LOCK_FILE = "/tmp/gitea_issue_lock.json" ISSUE_LOCK_FILE = "/tmp/gitea_issue_lock.json"
WORK_LEASE_TTL_HOURS = 4
AUTHOR_ISSUE_WORK_LEASE = "author_issue_work" AUTHOR_ISSUE_WORK_LEASE = "author_issue_work"
VALID_WORK_LEASE_OPERATIONS = frozenset({ VALID_WORK_LEASE_OPERATIONS = frozenset({
AUTHOR_ISSUE_WORK_LEASE, AUTHOR_ISSUE_WORK_LEASE,
@@ -2562,12 +2562,7 @@ def _build_author_issue_work_lease(
host: str | None, host: str | None,
) -> dict: ) -> dict:
created = _work_lease_now() created = _work_lease_now()
# #790 Slice A: the window comes from the central policy, not a literal here. expires = created + timedelta(hours=WORK_LEASE_TTL_HOURS)
# 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 { return {
"operation_type": AUTHOR_ISSUE_WORK_LEASE, "operation_type": AUTHOR_ISSUE_WORK_LEASE,
"issue_number": issue_number, "issue_number": issue_number,
@@ -2578,15 +2573,6 @@ def _build_author_issue_work_lease(
"created_at": _work_lease_timestamp(created), "created_at": _work_lease_timestamp(created),
"expires_at": _work_lease_timestamp(expires), "expires_at": _work_lease_timestamp(expires),
"last_heartbeat_at": _work_lease_timestamp(created), "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,
} }
@@ -4356,138 +4342,6 @@ def gitea_lock_issue(
return result 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() @mcp.tool()
def gitea_assess_work_issue_duplicate( def gitea_assess_work_issue_duplicate(
issue_number: int, issue_number: int,
+29 -553
View File
@@ -15,27 +15,15 @@ import json
import os import os
import re import re
import tempfile import tempfile
import uuid
from contextlib import contextmanager from contextlib import contextmanager
from datetime import datetime, timedelta, timezone from datetime import datetime, timedelta, timezone
from typing import Any from typing import Any
import lease_policy
LOCK_DIR_ENV = "GITEA_ISSUE_LOCK_DIR" LOCK_DIR_ENV = "GITEA_ISSUE_LOCK_DIR"
DEFAULT_LOCK_DIR = os.path.expanduser("~/.cache/gitea-tools/issue-locks") DEFAULT_LOCK_DIR = os.path.expanduser("~/.cache/gitea-tools/issue-locks")
WORK_LEASE_TTL_HOURS = 4
AUTHOR_ISSUE_WORK_LEASE = "author_issue_work" 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._+-]+") _SAFE_SEGMENT_RE = re.compile(r"[^A-Za-z0-9._+-]+")
@@ -265,331 +253,6 @@ def bind_session_lock(
return path 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: def read_session_issue_lock(lock_dir: str | None = None) -> dict[str, Any] | None:
root = (lock_dir or default_lock_dir()).strip() root = (lock_dir or default_lock_dir()).strip()
pointer = read_lock_file(session_pointer_path(root)) pointer = read_lock_file(session_pointer_path(root))
@@ -673,16 +336,6 @@ def _parse_lease_timestamp(value: str | None) -> datetime | None:
return 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: def lease_expires_at(lock: dict[str, Any] | None) -> datetime | None:
if not lock: if not lock:
return None return None
@@ -703,216 +356,60 @@ def is_lease_live(lock: dict[str, Any] | None, *, now: datetime | None = None) -
return assess_lock_freshness(lock, now=now)["live"] 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( def assess_lock_freshness(
lock_data: dict[str, Any] | None, lock_data: dict[str, Any] | None,
*, *,
now: datetime | None = None, now: datetime | None = None,
) -> dict[str, Any]: ) -> 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) current = _lease_now(now)
if not lock_data: if not lock_data:
return { return {
"status": STATUS_ABSENT, "status": "absent",
"live": False, "live": False,
"stale": False, "stale": False,
"reason": "no lock record", "reason": "no lock record",
} }
lease = lock_data.get("work_lease")
expires_at = lease_expires_at(lock_data) expires_at = lease_expires_at(lock_data)
heartbeat_at = _lease_heartbeat_at(lock_data) lease = lock_data.get("work_lease")
created_at = ( heartbeat_at = _parse_lease_timestamp(lock_data.get("last_heartbeat_at"))
_parse_lease_timestamp(lease.get("created_at")) if heartbeat_at is None and isinstance(lease, dict):
if isinstance(lease, dict) heartbeat_at = _parse_lease_timestamp(lease.get("last_heartbeat_at"))
else None
)
pid = lock_data.get("session_pid") pid = lock_data.get("session_pid")
if pid is None: if pid is None:
pid = lock_data.get("pid") pid = lock_data.get("pid")
# Evidence only. Never consulted to grant liveness (AC-N2).
pid_alive = is_process_alive(pid) if pid is not None else False pid_alive = is_process_alive(pid) if pid is not None else False
lifecycle = lease_lifecycle_version(lock_data) if expires_at and expires_at <= current:
legacy = lifecycle != lease_policy.LIFECYCLE_HEARTBEAT_V1 return {
policy = lease_policy.policy_for(lease_task_class(lock_data)) "status": "expired",
"live": False,
"stale": True,
"reason": f"lease expired at {expires_at.isoformat()}",
"pid_alive": pid_alive,
}
evidence: dict[str, Any] = { 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",
"pid_alive": pid_alive, "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, "heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None,
"expires_at": expires_at.isoformat() if expires_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: def _same_realpath(left: str | None, right: str | None) -> bool:
if not left or not right: if not left or not right:
@@ -949,27 +446,6 @@ def assess_expired_lock_reclaim(
"reasons": ["lock is still live; cannot reclaim (fail closed)"], "reasons": ["lock is still live; cannot reclaim (fail closed)"],
"freshness": freshness, "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") pid = existing_lock.get("session_pid")
if pid is None: if pid is None:
pid = existing_lock.get("pid") pid = existing_lock.get("pid")
-212
View File
@@ -1,212 +0,0 @@
"""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,15 +32,6 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.issue.comment", "permission": "gitea.issue.comment",
"role": "author", "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": { "set_issue_labels": {
"permission": "gitea.issue.comment", "permission": "gitea.issue.comment",
"role": "author", "role": "author",
-444
View File
@@ -1,444 +0,0 @@
"""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
@@ -1,594 +0,0 @@
"""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()
+499
View File
@@ -0,0 +1,499 @@
"""Tests for the read-only system-health API (#634).
Covers the acceptance criteria directly: a structured payload with readiness
and a dependency list (AC1), version and uptime when knowable (AC2), stale
runtime reported without a false mutation-safe claim (AC3), and the healthy /
degraded-dependency / redaction cases (AC4).
"""
import json
import os
import sqlite3
import sys
import tempfile
import unittest
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient
import control_plane_db
from webui.app import create_app
from webui.deployment_boundary import scan_text_for_client_secrets
from webui.system_health import (
API_PATH,
STATUS_DEGRADED,
STATUS_DOWN,
STATUS_OK,
STATUS_SKIPPED,
DependencyProbe,
StaleRuntime,
assess_stale_runtime,
clear_probe_cache,
load_system_health,
namespace_summaries,
probe_control_plane_db,
probe_gitea,
process_uptime,
redact,
redact_url,
snapshot_to_dict,
)
def _probe(name, status, *, required=True, detail="detail", kind="test"):
return DependencyProbe(
name=name,
kind=kind,
status=status,
detail=detail,
required=required,
latency_ms=1.5,
metadata={},
)
_ALL_HEALTHY = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_OK, required=False, kind="http"),
)
_CLEAN_PARITY = StaleRuntime(
daemon_head="abc123",
checkout_head="abc123",
remote_head="abc123",
stale=False,
determinable=True,
mutation_safe=True,
reasons=(),
)
class CleanParityMixin:
"""Pin parity for tests about aggregation rather than staleness.
Without this the assertions depend on the real checkout: a worktree whose
branch is ahead of its upstream is genuinely stale, which would degrade the
overall status and make these cases fail for an unrelated reason.
"""
def setUp(self):
super().setUp()
patcher = mock.patch(
"webui.system_health.assess_stale_runtime",
return_value=_CLEAN_PARITY,
)
patcher.start()
self.addCleanup(patcher.stop)
class TestDependencyAggregation(CleanParityMixin, unittest.TestCase):
"""AC1 — readiness and dependency list derived from probe results."""
def test_all_healthy_is_ok_and_ready(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc123")
self.assertEqual(snapshot.status, STATUS_OK)
self.assertTrue(snapshot.ready)
self.assertTrue(snapshot.readiness_complete)
self.assertEqual(snapshot.readiness_reasons, ())
self.assertEqual(len(snapshot.dependencies), 3)
def test_required_dependency_down_blocks_readiness(self):
probes = (
_probe("control_plane_db", STATUS_DOWN, detail="file missing", kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_OK, required=False, kind="http"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertEqual(snapshot.status, STATUS_DOWN)
self.assertFalse(snapshot.ready)
self.assertTrue(
any("control_plane_db" in reason for reason in snapshot.readiness_reasons)
)
def test_optional_dependency_down_degrades_but_stays_ready(self):
"""A failing optional probe must not claim the process itself is unready."""
probes = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_DOWN, required=False, detail="timeout", kind="http"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertEqual(snapshot.status, STATUS_DEGRADED)
self.assertTrue(snapshot.ready)
self.assertTrue(any("gitea" in reason for reason in snapshot.readiness_reasons))
def test_unrun_required_probe_leaves_readiness_incomplete(self):
"""Not probed is not the same as passing."""
probes = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_SKIPPED, detail="offline", kind="git"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertFalse(snapshot.ready)
self.assertFalse(snapshot.readiness_complete)
self.assertEqual(snapshot.status, STATUS_DEGRADED)
def test_skipped_optional_probe_does_not_block_readiness(self):
probes = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_SKIPPED, required=False, kind="http"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertTrue(snapshot.ready)
self.assertTrue(snapshot.readiness_complete)
class TestVersionAndUptime(CleanParityMixin, unittest.TestCase):
"""AC2 — version and uptime present when knowable."""
def test_uptime_and_start_time_present(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc123")
self.assertGreaterEqual(snapshot.uptime_seconds, 0.0)
self.assertIn("T", snapshot.started_at)
def test_process_uptime_helper_matches_shape(self):
started_at, uptime = process_uptime()
self.assertIn("T", started_at)
self.assertGreaterEqual(uptime, 0.0)
def test_version_reports_python_and_schema_version(self):
probes = (
DependencyProbe(
name="control_plane_db",
kind="sqlite",
status=STATUS_OK,
detail="ok",
required=True,
latency_ms=1.0,
metadata={"schema_version": control_plane_db.SCHEMA_VERSION},
),
_probe("repository", STATUS_OK, kind="git"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertEqual(
snapshot.version.control_plane_schema_version,
control_plane_db.SCHEMA_VERSION,
)
self.assertTrue(snapshot.version.python_version)
def test_version_known_flag_false_when_sha_unavailable(self):
with mock.patch("webui.system_health._git", return_value=None):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc")
self.assertIsNone(snapshot.version.git_sha)
self.assertFalse(snapshot.version.known)
class TestStaleRuntime(unittest.TestCase):
"""AC3 — stale runtime reflected without a false mutation-safe claim."""
def test_matching_commits_are_mutation_safe(self):
assessment = assess_stale_runtime(
Path("/tmp"),
daemon_head="aaa",
git_reader=lambda *args: "aaa",
)
self.assertFalse(assessment.stale)
self.assertTrue(assessment.determinable)
self.assertTrue(assessment.mutation_safe)
def test_diverged_commits_are_stale_and_not_mutation_safe(self):
reads = {"HEAD": "aaa", "@{upstream}": "bbb"}
assessment = assess_stale_runtime(
Path("/tmp"),
daemon_head="aaa",
git_reader=lambda *args: reads.get(args[-1]),
)
self.assertTrue(assessment.stale)
self.assertFalse(assessment.mutation_safe)
self.assertTrue(assessment.reasons)
def test_unknown_remote_is_not_mutation_safe(self):
"""Indeterminate must never read as safe."""
reads = {"HEAD": "aaa", "@{upstream}": None}
assessment = assess_stale_runtime(
Path("/tmp"),
daemon_head="aaa",
git_reader=lambda *args: reads.get(args[-1]),
)
self.assertFalse(assessment.determinable)
self.assertFalse(assessment.mutation_safe)
self.assertFalse(assessment.stale)
self.assertTrue(
any("indeterminate" in reason for reason in assessment.reasons)
)
def test_unobservable_daemon_head_is_disclosed(self):
assessment = assess_stale_runtime(
Path("/tmp"),
git_reader=lambda *args: "aaa",
)
self.assertTrue(
any("not observable" in reason for reason in assessment.reasons)
)
def test_stale_runtime_degrades_overall_status(self):
reads = {"HEAD": "aaa", "@{upstream}": "bbb"}
# Pinned rather than inherited: this path uses the default git reader,
# so the assertion must hold whether or not the suite runs offline.
with mock.patch.dict(os.environ, {"WEBUI_TEST_OFFLINE": ""}), mock.patch(
"webui.system_health._git",
side_effect=lambda repo, *args: reads.get(args[-1]),
):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="aaa")
self.assertTrue(snapshot.stale_runtime.stale)
self.assertFalse(snapshot.stale_runtime.mutation_safe)
self.assertEqual(snapshot.status, STATUS_DEGRADED)
class TestControlPlaneDbProbe(unittest.TestCase):
"""The required local dependency, probed read-only."""
def setUp(self):
self.tmp = tempfile.TemporaryDirectory()
self.addCleanup(self.tmp.cleanup)
self.db_path = str(Path(self.tmp.name) / "control-plane.db")
def _build_db(self, schema_version):
conn = sqlite3.connect(self.db_path)
conn.execute("CREATE TABLE schema_meta (key TEXT PRIMARY KEY, value TEXT)")
conn.execute("CREATE TABLE leases (lease_id TEXT PRIMARY KEY, status TEXT)")
conn.execute(
"INSERT INTO schema_meta(key, value) VALUES ('schema_version', ?)",
(str(schema_version),),
)
conn.execute("INSERT INTO leases(lease_id, status) VALUES ('l1', 'active')")
conn.commit()
conn.close()
def test_missing_database_is_down(self):
probe = probe_control_plane_db(str(Path(self.tmp.name) / "absent.db"))
self.assertEqual(probe.status, STATUS_DOWN)
self.assertTrue(probe.required)
self.assertIsNotNone(probe.latency_ms)
def test_matching_schema_is_ok(self):
self._build_db(control_plane_db.SCHEMA_VERSION)
probe = probe_control_plane_db(self.db_path)
self.assertEqual(probe.status, STATUS_OK)
self.assertEqual(
probe.metadata["schema_version"], control_plane_db.SCHEMA_VERSION
)
self.assertEqual(probe.metadata["active_leases"], 1)
def test_mismatched_schema_is_degraded(self):
self._build_db(control_plane_db.SCHEMA_VERSION + 99)
probe = probe_control_plane_db(self.db_path)
self.assertEqual(probe.status, STATUS_DEGRADED)
def test_probe_does_not_create_a_database(self):
"""A health check must never initialise the substrate it inspects."""
absent = str(Path(self.tmp.name) / "never-created.db")
probe_control_plane_db(absent)
self.assertFalse(Path(absent).exists())
def test_unreadable_database_is_down_not_raised(self):
Path(self.db_path).write_text("this is not a sqlite database")
probe = probe_control_plane_db(self.db_path)
self.assertEqual(probe.status, STATUS_DOWN)
class TestRedaction(unittest.TestCase):
"""AC4 — redaction. No credential-shaped text crosses the boundary."""
def test_redacts_token_assignment(self):
cleaned = redact("failed with token=ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ012345")
self.assertNotIn("ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ012345", cleaned)
self.assertIn("[redacted]", cleaned)
def test_redacts_authorization_header_text(self):
cleaned = redact("Authorization: Bearer abcdefghijklmnopqrstuvwxyz123456")
self.assertNotIn("abcdefghijklmnopqrstuvwxyz123456", cleaned)
def test_redacts_long_opaque_strings(self):
cleaned = redact("value 0123456789abcdef0123456789abcdef here")
self.assertNotIn("0123456789abcdef0123456789abcdef", cleaned)
def test_url_userinfo_and_query_are_stripped(self):
cleaned = redact_url("https://user:[email protected]/api/v1?token=xyz")
self.assertNotIn("secretpass", cleaned)
self.assertNotIn("token=xyz", cleaned)
self.assertEqual(cleaned, "https://gitea.example.com/api/v1")
def test_url_inside_free_text_is_redacted(self):
cleaned = redact("GET https://u:[email protected]/x?token=abc failed")
self.assertNotIn("u:p@", cleaned)
self.assertNotIn("token=abc", cleaned)
def test_gitea_probe_failure_detail_is_redacted(self):
boom = RuntimeError(
"connection refused for https://user:[email protected]/api/v1/version"
)
with mock.patch("webui.system_health.get_auth_header", return_value="token x"), \
mock.patch("webui.system_health.api_request", side_effect=boom):
probe = probe_gitea("gitea.example.com")
self.assertEqual(probe.status, STATUS_DOWN)
self.assertNotIn("hunter2", probe.detail)
self.assertEqual(scan_text_for_client_secrets(probe.detail), [])
def test_credential_guard_refusal_is_a_status_not_a_crash(self):
with mock.patch(
"webui.system_health.get_auth_header",
side_effect=RuntimeError("daemon guard refused"),
):
probe = probe_gitea("gitea.example.com")
self.assertEqual(probe.status, STATUS_DEGRADED)
self.assertFalse(probe.required)
class TestNamespaceSummaries(unittest.TestCase):
"""A web process cannot prove IDE namespace health, and must not claim to."""
def test_every_namespace_reports_unproven(self):
rows = namespace_summaries()
self.assertTrue(rows)
for row in rows:
with self.subTest(namespace=row["namespace"]):
self.assertEqual(row["status"], "unproven")
self.assertFalse(row["ide_namespace_proven"])
self.assertIn("client_namespace", row["reason"])
class TestSystemHealthRoutes(CleanParityMixin, unittest.TestCase):
"""The HTTP surface: versioned path, status codes, read-only guard."""
def setUp(self):
super().setUp()
clear_probe_cache()
self.addCleanup(clear_probe_cache)
self.client = TestClient(create_app())
def _patch_snapshot(self, probes, daemon_head="abc123"):
snapshot = load_system_health(probes=probes, daemon_head=daemon_head)
patcher = mock.patch(
"webui.app.load_system_health",
return_value=snapshot,
)
patcher.start()
self.addCleanup(patcher.stop)
return snapshot
def test_versioned_route_is_registered(self):
self.assertEqual(API_PATH, "/api/v1/system/health")
self._patch_snapshot(_ALL_HEALTHY)
response = self.client.get(API_PATH)
self.assertEqual(response.status_code, 200)
def test_healthy_payload_shape(self):
self._patch_snapshot(_ALL_HEALTHY)
data = self.client.get(API_PATH).json()
self.assertEqual(data["status"], STATUS_OK)
self.assertTrue(data["readiness"]["ready"])
self.assertTrue(data["readiness"]["complete"])
self.assertEqual(data["api"], API_PATH)
self.assertEqual(len(data["dependencies"]), 3)
for key in ("version", "process", "stale_runtime", "mcp_namespaces"):
self.assertIn(key, data)
self.assertIn("uptime_seconds", data["process"])
self.assertIn("mutation_safe", data["stale_runtime"])
def test_degraded_dependency_returns_503(self):
probes = (
_probe("control_plane_db", STATUS_DOWN, detail="missing", kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
)
self._patch_snapshot(probes)
response = self.client.get(API_PATH)
self.assertEqual(response.status_code, 503)
data = response.json()
self.assertFalse(data["readiness"]["ready"])
self.assertTrue(data["readiness"]["reasons"])
def test_dependency_entries_expose_status_and_latency(self):
self._patch_snapshot(_ALL_HEALTHY)
data = self.client.get(API_PATH).json()
names = {entry["name"] for entry in data["dependencies"]}
self.assertEqual(names, {"control_plane_db", "repository", "gitea"})
for entry in data["dependencies"]:
with self.subTest(dependency=entry["name"]):
self.assertIn("status", entry)
self.assertIn("required", entry)
self.assertIn("latency_ms", entry)
def test_response_body_carries_no_client_secrets(self):
self._patch_snapshot(_ALL_HEALTHY)
body = self.client.get(API_PATH).text
self.assertEqual(scan_text_for_client_secrets(body), [])
def test_deep_flag_is_forwarded(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc")
with mock.patch(
"webui.app.load_system_health", return_value=snapshot
) as loader:
self.client.get(f"{API_PATH}?deep=1")
loader.assert_called_once_with(deep=True)
def test_shallow_is_the_default(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc")
with mock.patch(
"webui.app.load_system_health", return_value=snapshot
) as loader:
self.client.get(API_PATH)
loader.assert_called_once_with(deep=False)
def test_route_rejects_mutation_methods(self):
for method in ("POST", "PUT", "PATCH", "DELETE"):
with self.subTest(method=method):
response = self.client.request(method, API_PATH)
self.assertEqual(response.status_code, 405)
self.assertEqual(response.json()["error"], "read-only-mvp")
def test_default_shallow_call_skips_the_network_probe(self):
"""The expensive probe must not run unless it was asked for."""
with mock.patch("webui.system_health.probe_gitea") as probe:
snapshot = load_system_health(deep=False)
probe.assert_not_called()
gitea = next(p for p in snapshot.dependencies if p.name == "gitea")
self.assertEqual(gitea.status, STATUS_SKIPPED)
class TestHealthRouteBackwardCompatibility(unittest.TestCase):
"""`/health` is expanded additively; MVP consumers must keep working."""
def setUp(self):
self.client = TestClient(create_app())
def test_mvp_keys_are_unchanged(self):
data = self.client.get("/health").json()
self.assertEqual(data["status"], "ok")
self.assertEqual(data["service"], "mcp-control-plane-webui")
self.assertEqual(data["mode"], "read-only-mvp")
self.assertIn("timestamp", data)
self.assertEqual(data["deployment"]["mode"], "internal-operator-console")
def test_health_points_at_the_versioned_api(self):
data = self.client.get("/health").json()
self.assertEqual(data["system_health_api"], API_PATH)
self.assertIn("uptime_seconds", data)
self.assertIn("started_at", data)
def test_health_runs_no_dependency_probe(self):
"""Liveness must stay cheap: no probe, no snapshot assembly."""
with mock.patch("webui.app.load_system_health") as loader:
response = self.client.get("/health")
self.assertEqual(response.status_code, 200)
loader.assert_not_called()
class TestSnapshotSerialisation(CleanParityMixin, unittest.TestCase):
def test_snapshot_dict_is_json_serialisable(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc123")
encoded = json.dumps(snapshot_to_dict(snapshot))
self.assertIn("readiness", encoded)
if __name__ == "__main__":
unittest.main()
+34
View File
@@ -45,6 +45,12 @@ from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as wo
from webui.worktree_views import render_worktrees_page from webui.worktree_views import render_worktrees_page
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
from webui.runtime_views import render_runtime_page from webui.runtime_views import render_runtime_page
from webui.system_health import (
API_PATH as SYSTEM_HEALTH_API_PATH,
load_system_health,
process_uptime,
snapshot_to_dict as system_health_to_dict,
)
_READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"}) _READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"})
_AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"}) _AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"})
@@ -78,16 +84,43 @@ async def home(_request: Request) -> HTMLResponse:
async def health(_request: Request) -> JSONResponse: async def health(_request: Request) -> JSONResponse:
"""Liveness only — deliberately cheap, runs no dependency probe (#634).
Every MVP key is retained so existing pollers keep working; the additions
are a pointer to the structured API and the in-memory process uptime.
Readiness lives at that API because answering it costs real probes.
"""
bind_host = _request.app.state.webui_bind_host bind_host = _request.app.state.webui_bind_host
started_at, uptime_seconds = process_uptime()
return JSONResponse({ return JSONResponse({
"status": "ok", "status": "ok",
"service": "mcp-control-plane-webui", "service": "mcp-control-plane-webui",
"mode": "read-only-mvp", "mode": "read-only-mvp",
"timestamp": datetime.now(timezone.utc).isoformat(), "timestamp": datetime.now(timezone.utc).isoformat(),
"deployment": deployment_snapshot(bind_host=bind_host), "deployment": deployment_snapshot(bind_host=bind_host),
"started_at": started_at,
"uptime_seconds": uptime_seconds,
"system_health_api": SYSTEM_HEALTH_API_PATH,
}) })
def _truthy_flag(value: str | None) -> bool:
return (value or "").strip().lower() in {"1", "true", "yes", "on"}
async def api_system_health(request: Request) -> JSONResponse:
"""Structured read-only system health (#634).
`?deep=1` opts into the expensive network probe. The response status code
reflects readiness so automated checks can branch on it without parsing the
body: 200 when ready, 503 when a required dependency failed or never ran.
"""
deep = _truthy_flag(request.query_params.get("deep"))
snapshot = load_system_health(deep=deep)
payload = system_health_to_dict(snapshot)
return JSONResponse(payload, status_code=200 if snapshot.ready else 503)
async def queue(_request: Request) -> HTMLResponse: async def queue(_request: Request) -> HTMLResponse:
snapshot = load_queue_snapshot() snapshot = load_queue_snapshot()
return HTMLResponse(render_page(title="Queue", body_html=render_queue_page(snapshot))) return HTMLResponse(render_page(title="Queue", body_html=render_queue_page(snapshot)))
@@ -399,6 +432,7 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
routes=[ routes=[
Route("/", home, methods=["GET"]), Route("/", home, methods=["GET"]),
Route("/health", health, methods=["GET"]), Route("/health", health, methods=["GET"]),
Route(SYSTEM_HEALTH_API_PATH, api_system_health, methods=["GET"]),
Route("/queue", queue, methods=["GET"]), Route("/queue", queue, methods=["GET"]),
Route("/api/queue", api_queue, methods=["GET"]), Route("/api/queue", api_queue, methods=["GET"]),
Route("/projects", projects, methods=["GET"]), Route("/projects", projects, methods=["GET"]),
+682
View File
@@ -0,0 +1,682 @@
"""Read-only system-health model for the operator console API (#634).
`/health` answers liveness only. Operators automating readiness checks need a
structured view of *why* the control plane is or is not usable: which
dependencies answered, how long they took, what version of the code is running,
and whether the runtime is stale relative to its remote.
Three rules shape this module.
* **Read-only.** Every probe opens its subject read-only. The control-plane
database is opened through a ``mode=ro`` URI so a health check can never
create or migrate a schema, and no probe writes, restarts, or reloads
anything — restart controls are Phase 2, and #630 forbids process-kill
recovery outright.
* **Fail-soft.** A dependency that is unreachable is a *status*, not an
exception. Probes catch their own failures and report them as a degraded or
down entry carrying a reason.
* **Never claim more than was proven.** Readiness is derived only from probes
that actually ran, ``mutation_safe`` stays false unless the parity commits are
known and equal, and an MCP namespace is reported unproven because a web
process cannot exercise the IDE-managed client path (#543).
"""
from __future__ import annotations
import os
import re
import sqlite3
import subprocess
import time
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Callable
from urllib.parse import urlsplit, urlunsplit
import control_plane_db
import mcp_namespace_health
from gitea_auth import api_request, get_auth_header, gitea_url
from webui.project_registry import load_registry
SERVICE_NAME = "mcp-control-plane-webui"
API_PATH = "/api/v1/system/health"
STATUS_OK = "ok"
STATUS_DEGRADED = "degraded"
STATUS_DOWN = "down"
STATUS_SKIPPED = "skipped"
STATUS_UNPROVEN = "unproven"
# Statuses that count as a healthy answer from a probe.
_HEALTHY_STATUSES = frozenset({STATUS_OK})
# Statuses meaning "this probe did not run", as opposed to "it ran and failed".
_NOT_RUN_STATUSES = frozenset({STATUS_SKIPPED})
_DEEP_PROBE_TTL_ENV = "WEBUI_HEALTH_PROBE_TTL_SECONDS"
_DEFAULT_DEEP_PROBE_TTL = 15.0
_GITEA_PROBE_TIMEOUT_SECONDS = 5.0
_OFFLINE_ENV = "WEBUI_TEST_OFFLINE"
# Credential-shaped material that must never reach the browser, mirroring the
# forbidden client patterns in webui/deployment_boundary.py.
_SECRET_RE = re.compile(
r"(?i)\b(token|password|passwd|secret|authorization|bearer)\b\s*[:=]?\s*\S+"
)
_LONG_OPAQUE_RE = re.compile(r"\b[A-Za-z0-9_\-]{32,}\b")
# Captured once at import so uptime measures this process, not the request.
_STARTED_AT = datetime.now(timezone.utc)
_STARTED_MONOTONIC = time.monotonic()
# TTL cache for the expensive (network) probe only.
_deep_cache: dict[str, tuple[float, "DependencyProbe"]] = {}
@dataclass(frozen=True)
class DependencyProbe:
"""One dependency check, fail-soft, with its own latency."""
name: str
kind: str
status: str
detail: str
required: bool
latency_ms: float | None = None
metadata: dict[str, Any] | None = None
@property
def healthy(self) -> bool:
return self.status in _HEALTHY_STATUSES
@property
def ran(self) -> bool:
return self.status not in _NOT_RUN_STATUSES
@dataclass(frozen=True)
class VersionInfo:
git_sha: str | None
git_describe: str | None
control_plane_schema_version: int | None
python_version: str
known: bool
@dataclass(frozen=True)
class StaleRuntime:
"""Parity between the running code, the checkout, and the remote.
``mutation_safe`` is deliberately conservative: unknown is not safe.
"""
daemon_head: str | None
checkout_head: str | None
remote_head: str | None
stale: bool
determinable: bool
mutation_safe: bool
reasons: tuple[str, ...]
@dataclass(frozen=True)
class SystemHealthSnapshot:
status: str
ready: bool
readiness_complete: bool
readiness_reasons: tuple[str, ...]
service: str
mode: str
version: VersionInfo
started_at: str
uptime_seconds: float
timestamp: str
deep_probes_requested: bool
dependencies: tuple[DependencyProbe, ...]
mcp_namespaces: tuple[dict[str, Any], ...]
stale_runtime: StaleRuntime
probe_errors: tuple[str, ...] = ()
def process_uptime() -> tuple[str, float]:
"""Process start timestamp and uptime — in-memory, safe for `/health`."""
return _STARTED_AT.isoformat(), round(time.monotonic() - _STARTED_MONOTONIC, 3)
def _offline() -> bool:
return (os.environ.get(_OFFLINE_ENV) or "").strip().lower() in {"1", "true", "yes"}
def _repo_root() -> Path:
override = (os.environ.get("WEBUI_REPO_ROOT") or "").strip()
if override:
return Path(override).resolve()
return Path(__file__).resolve().parent.parent
def _deep_probe_ttl() -> float:
raw = (os.environ.get(_DEEP_PROBE_TTL_ENV) or "").strip()
if not raw:
return _DEFAULT_DEEP_PROBE_TTL
try:
value = float(raw)
except ValueError:
return _DEFAULT_DEEP_PROBE_TTL
return value if value >= 0 else _DEFAULT_DEEP_PROBE_TTL
def redact(text: str) -> str:
"""Strip credential-shaped material from operator-visible probe text.
Probe details carry exception strings, and an exception raised by an HTTP
client can quote the request that failed. Redaction happens here, at the
boundary where those strings become part of a browser-bound payload.
"""
if not text:
return ""
cleaned = _redact_urls(text)
cleaned = _SECRET_RE.sub(lambda m: f"{m.group(1)}=[redacted]", cleaned)
return _LONG_OPAQUE_RE.sub("[redacted]", cleaned)
def _redact_urls(text: str) -> str:
return re.sub(r"https?://\S+", lambda m: redact_url(m.group(0)), text)
def redact_url(url: str) -> str:
"""Reduce a URL to scheme://host/path — no userinfo, no query, no fragment."""
try:
parts = urlsplit(url)
except ValueError:
return "[redacted-url]"
if not parts.scheme or not parts.hostname:
return "[redacted-url]"
netloc = parts.hostname
if parts.port:
netloc = f"{netloc}:{parts.port}"
return urlunsplit((parts.scheme, netloc, parts.path, "", ""))
def _git(repo: Path, *args: str) -> str | None:
try:
completed = subprocess.run(
["git", "-C", str(repo), *args],
capture_output=True,
text=True,
check=False,
timeout=10,
)
except (OSError, subprocess.SubprocessError):
return None
if completed.returncode != 0:
return None
return (completed.stdout or "").strip() or None
def _load_version(repo: Path, *, schema_version: int | None) -> VersionInfo:
import platform
git_sha = None if _offline() else _git(repo, "rev-parse", "HEAD")
describe = None if _offline() else _git(repo, "describe", "--tags", "--always")
return VersionInfo(
git_sha=git_sha,
git_describe=describe,
control_plane_schema_version=schema_version,
python_version=platform.python_version(),
known=bool(git_sha),
)
# ---------------------------------------------------------------------------
# Dependency probes
# ---------------------------------------------------------------------------
def _elapsed_ms(started: float) -> float:
return round((time.monotonic() - started) * 1000, 3)
def probe_control_plane_db(db_path: str | None = None) -> DependencyProbe:
"""Read-only reachability check for the control-plane SQLite substrate.
Opened through a ``mode=ro`` URI on purpose: ``ControlPlaneDB.__init__``
creates directories and runs schema migrations, which a health check must
never do.
"""
path = (db_path or control_plane_db.default_db_path()).strip()
started = time.monotonic()
metadata: dict[str, Any] = {"path": path}
def _result(status: str, detail: str) -> DependencyProbe:
return DependencyProbe(
name="control_plane_db",
kind="sqlite",
status=status,
detail=detail,
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
if not path or not os.path.exists(path):
return _result(STATUS_DOWN, "control-plane database file does not exist yet")
try:
conn = sqlite3.connect(f"file:{path}?mode=ro", uri=True, timeout=5)
try:
row = conn.execute(
"SELECT value FROM schema_meta WHERE key = 'schema_version'"
).fetchone()
leases = conn.execute(
"SELECT COUNT(*) FROM leases WHERE status = 'active'"
).fetchone()
finally:
conn.close()
except sqlite3.Error as exc:
return _result(STATUS_DOWN, redact(f"control-plane database unreadable: {exc}"))
schema_version = int(row[0]) if row and str(row[0]).isdigit() else None
metadata["schema_version"] = schema_version
metadata["active_leases"] = int(leases[0]) if leases else None
if schema_version is None:
return _result(
STATUS_DEGRADED, "control-plane database has no recorded schema version"
)
if schema_version != control_plane_db.SCHEMA_VERSION:
return _result(
STATUS_DEGRADED,
f"control-plane schema version {schema_version} does not match the "
f"version this code expects ({control_plane_db.SCHEMA_VERSION})",
)
return _result(STATUS_OK, f"schema v{schema_version} readable")
def probe_repository(repo: Path) -> DependencyProbe:
"""Local checkout reachability — required, cheap, no network."""
started = time.monotonic()
metadata: dict[str, Any] = {"repo_root": str(repo)}
if _offline():
return DependencyProbe(
name="repository",
kind="git",
status=STATUS_SKIPPED,
detail=f"{_OFFLINE_ENV} is set; git probe skipped",
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
head = _git(repo, "rev-parse", "HEAD")
if not head:
return DependencyProbe(
name="repository",
kind="git",
status=STATUS_DOWN,
detail=f"HEAD could not be read at {repo}",
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
branch = _git(repo, "rev-parse", "--abbrev-ref", "HEAD")
metadata["head"] = head
metadata["branch"] = branch
return DependencyProbe(
name="repository",
kind="git",
status=STATUS_OK,
detail=f"checkout readable at {branch or 'detached HEAD'}",
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
def probe_gitea(host: str) -> DependencyProbe:
"""Live Gitea reachability. Expensive (network), so opt-in via ``deep``.
Optional by design: the console stays useful for local inventory when the
remote is unreachable, so a failure here degrades status without claiming
the process itself is unready.
"""
started = time.monotonic()
metadata: dict[str, Any] = {"host": host}
def _failure(status: str, detail: str) -> DependencyProbe:
return DependencyProbe(
name="gitea",
kind="http",
status=status,
detail=detail,
required=False,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
if not host:
return _failure(STATUS_DEGRADED, "no Gitea host is configured in the registry")
try:
auth = get_auth_header(host)
except Exception as exc: # noqa: BLE001 — credential guards are a status here
return _failure(STATUS_DEGRADED, redact(f"credential lookup refused: {exc}"))
if not auth:
return _failure(STATUS_DEGRADED, f"no credentials available for {host}")
url = gitea_url(host, "/api/v1/version")
metadata["endpoint"] = redact_url(url)
try:
data = api_request("GET", url, auth, timeout=_GITEA_PROBE_TIMEOUT_SECONDS)
except Exception as exc: # noqa: BLE001 — a down dependency is a status
return _failure(STATUS_DOWN, redact(f"Gitea probe failed: {exc}"))
if isinstance(data, dict) and data.get("version"):
metadata["gitea_version"] = str(data["version"])
return DependencyProbe(
name="gitea",
kind="http",
status=STATUS_OK,
detail=f"{host} reachable",
required=False,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
def _skipped_gitea(host: str) -> DependencyProbe:
return DependencyProbe(
name="gitea",
kind="http",
status=STATUS_SKIPPED,
detail="network probe not requested; call with ?deep=1 to run it",
required=False,
latency_ms=None,
metadata={"host": host},
)
def namespace_summaries() -> tuple[dict[str, Any], ...]:
"""Declared MCP namespaces, each honestly reported as unproven.
The web process runs outside the IDE-managed MCP client, so it cannot
invoke a namespace tool. Per #543 only a ``client_namespace`` probe proves
that path, and inventing a healthy verdict here is exactly the false claim
the mutation gates exist to prevent.
"""
rows: list[dict[str, Any]] = []
for namespace, required_tool in sorted(
mcp_namespace_health.REQUIRED_NAMESPACE_TOOLS.items()
):
classification = mcp_namespace_health.classify_namespace_probe(
namespace,
required_tool=required_tool,
probe_result=None,
probe_source=mcp_namespace_health.PROBE_SOURCE_UNKNOWN,
)
rows.append(
{
"namespace": namespace,
"required_tool": required_tool,
"status": STATUS_UNPROVEN,
"ide_namespace_proven": bool(classification.get("ide_namespace_proven")),
"reason": (
"the web console cannot invoke the IDE-managed MCP client; "
"namespace health must be proven with a client_namespace "
"probe (#543)"
),
"error_type": classification.get("error_type"),
}
)
return tuple(rows)
def assess_stale_runtime(
repo: Path,
*,
daemon_head: str | None = None,
git_reader: Callable[..., str | None] | None = None,
) -> StaleRuntime:
"""Three-way parity view: running code, local checkout, remote-tracking ref.
``mutation_safe`` requires all three to be known and equal. Anything less —
including "the remote ref was never fetched" — is reported as not safe with
a reason, so an operator never reads an unproven green.
"""
reader = git_reader or (lambda *args: _git(repo, *args))
reasons: list[str] = []
# The offline switch suppresses real subprocess calls; an explicitly
# injected reader is already a substitute for them and is always used.
offline = _offline() and git_reader is None
checkout_head = None if offline else reader("rev-parse", "HEAD")
remote_head = None if offline else reader("rev-parse", "@{upstream}")
if offline:
reasons.append(f"{_OFFLINE_ENV} is set; parity commits were not read")
else:
if checkout_head is None:
reasons.append("local checkout HEAD could not be read")
if remote_head is None:
reasons.append(
"no remote-tracking commit is known for the current branch; "
"remote staleness is indeterminate (no fetch is performed here)"
)
effective_daemon = daemon_head if daemon_head is not None else checkout_head
if daemon_head is None:
reasons.append(
"the running MCP daemon's startup commit is not observable from the "
"web process; the checkout commit is reported in its place"
)
determinable = bool(checkout_head and remote_head and effective_daemon)
stale = bool(
determinable and len({checkout_head, remote_head, effective_daemon}) > 1
)
if stale:
reasons.append(
"runtime, checkout, and remote commits disagree; restart the MCP "
"server after updating the checkout before trusting capability gates"
)
return StaleRuntime(
daemon_head=effective_daemon,
checkout_head=checkout_head,
remote_head=remote_head,
stale=stale,
determinable=determinable,
mutation_safe=bool(determinable and not stale),
reasons=tuple(reasons),
)
# ---------------------------------------------------------------------------
# Snapshot assembly
# ---------------------------------------------------------------------------
def _default_host() -> str:
registry = load_registry()
if not registry.projects:
return ""
raw = registry.projects[0].remote_host
parts = urlsplit(raw.strip())
return parts.netloc or raw.strip().rstrip("/")
def _aggregate(
probes: tuple[DependencyProbe, ...],
) -> tuple[str, bool, bool, tuple[str, ...]]:
"""Fold probe results into overall status and readiness.
Required probes drive readiness; optional probes can only degrade status.
A probe that did not run leaves readiness incomplete rather than passing.
"""
reasons: list[str] = []
required = [probe for probe in probes if probe.required]
unrun_required = [probe for probe in required if not probe.ran]
failed_required = [probe for probe in required if probe.ran and not probe.healthy]
failed_optional = [
probe
for probe in probes
if not probe.required and probe.ran and not probe.healthy
]
for probe in unrun_required:
reasons.append(
f"required dependency '{probe.name}' was not probed: {probe.detail}"
)
for probe in failed_required:
reasons.append(
f"required dependency '{probe.name}' is {probe.status}: {probe.detail}"
)
for probe in failed_optional:
reasons.append(
f"optional dependency '{probe.name}' is {probe.status}: {probe.detail}"
)
readiness_complete = not unrun_required
ready = readiness_complete and not failed_required
if any(probe.status == STATUS_DOWN for probe in failed_required):
status = STATUS_DOWN
elif failed_required or failed_optional or unrun_required:
status = STATUS_DEGRADED
else:
status = STATUS_OK
return status, ready, readiness_complete, tuple(reasons)
def load_system_health(
*,
deep: bool = False,
host: str | None = None,
probes: tuple[DependencyProbe, ...] | None = None,
daemon_head: str | None = None,
use_cache: bool = True,
) -> SystemHealthSnapshot:
"""Assemble the read-only system-health snapshot.
``deep=True`` adds the network probe against Gitea; its result is cached for
a short TTL so repeated dashboard polls do not amplify into remote load.
"""
repo = _repo_root()
probe_errors: list[str] = []
if probes is None:
collected: list[DependencyProbe] = []
for probe_fn in (
lambda: probe_control_plane_db(),
lambda: probe_repository(repo),
):
try:
collected.append(probe_fn())
except Exception as exc: # noqa: BLE001 — a probe must not 500 the API
probe_errors.append(redact(f"probe raised: {exc}"))
resolved_host = host if host is not None else _default_host()
if deep and not _offline():
collected.append(_cached_gitea_probe(resolved_host, use_cache=use_cache))
else:
collected.append(_skipped_gitea(resolved_host))
probes = tuple(collected)
status, ready, readiness_complete, reasons = _aggregate(probes)
stale = assess_stale_runtime(repo, daemon_head=daemon_head)
if stale.stale:
if status == STATUS_OK:
status = STATUS_DEGRADED
reasons = reasons + (
"runtime is stale relative to its remote-tracking commit",
)
db_probe = next((p for p in probes if p.name == "control_plane_db"), None)
schema_version = None
if db_probe and db_probe.metadata:
schema_version = db_probe.metadata.get("schema_version")
return SystemHealthSnapshot(
status=status,
ready=ready,
readiness_complete=readiness_complete,
readiness_reasons=reasons,
service=SERVICE_NAME,
mode="read-only",
version=_load_version(repo, schema_version=schema_version),
started_at=_STARTED_AT.isoformat(),
uptime_seconds=round(time.monotonic() - _STARTED_MONOTONIC, 3),
timestamp=datetime.now(timezone.utc).isoformat(),
deep_probes_requested=deep,
dependencies=probes,
mcp_namespaces=namespace_summaries(),
stale_runtime=stale,
probe_errors=tuple(probe_errors),
)
def _cached_gitea_probe(host: str, *, use_cache: bool = True) -> DependencyProbe:
ttl = _deep_probe_ttl()
now = time.monotonic()
if use_cache and ttl > 0:
cached = _deep_cache.get(host)
if cached and (now - cached[0]) < ttl:
return cached[1]
probe = probe_gitea(host)
if use_cache and ttl > 0:
_deep_cache[host] = (now, probe)
return probe
def clear_probe_cache() -> None:
"""Drop cached deep-probe results (tests and operator-forced refresh)."""
_deep_cache.clear()
def probe_to_dict(probe: DependencyProbe) -> dict[str, Any]:
return {
"name": probe.name,
"kind": probe.kind,
"status": probe.status,
"detail": probe.detail,
"required": probe.required,
"healthy": probe.healthy,
"latency_ms": probe.latency_ms,
"metadata": dict(probe.metadata or {}),
}
def snapshot_to_dict(snapshot: SystemHealthSnapshot) -> dict[str, Any]:
return {
"status": snapshot.status,
"service": snapshot.service,
"mode": snapshot.mode,
"api": API_PATH,
"timestamp": snapshot.timestamp,
"readiness": {
"ready": snapshot.ready,
"complete": snapshot.readiness_complete,
"reasons": list(snapshot.readiness_reasons),
},
"version": {
"git_sha": snapshot.version.git_sha,
"git_describe": snapshot.version.git_describe,
"control_plane_schema_version": (
snapshot.version.control_plane_schema_version
),
"python_version": snapshot.version.python_version,
"known": snapshot.version.known,
},
"process": {
"started_at": snapshot.started_at,
"uptime_seconds": snapshot.uptime_seconds,
},
"deep_probes_requested": snapshot.deep_probes_requested,
"dependencies": [probe_to_dict(probe) for probe in snapshot.dependencies],
"mcp_namespaces": [dict(row) for row in snapshot.mcp_namespaces],
"stale_runtime": {
"daemon_head": snapshot.stale_runtime.daemon_head,
"checkout_head": snapshot.stale_runtime.checkout_head,
"remote_head": snapshot.stale_runtime.remote_head,
"stale": snapshot.stale_runtime.stale,
"determinable": snapshot.stale_runtime.determinable,
"mutation_safe": snapshot.stale_runtime.mutation_safe,
"reasons": list(snapshot.stale_runtime.reasons),
},
"probe_errors": list(snapshot.probe_errors),
}