From 5494696227b84148b374412e9b38d58c1eccfca5 Mon Sep 17 00:00:00 2001 From: jcwalker3 Date: Wed, 22 Jul 2026 15:59:30 -0500 Subject: [PATCH 1/4] feat(webui): read-only system-health API (Closes #634) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) --- docs/webui-local-dev.md | 82 +++- tests/test_webui_system_health.py | 499 ++++++++++++++++++++++ webui/app.py | 34 ++ webui/system_health.py | 682 ++++++++++++++++++++++++++++++ 4 files changed, 1296 insertions(+), 1 deletion(-) create mode 100644 tests/test_webui_system_health.py create mode 100644 webui/system_health.py diff --git a/docs/webui-local-dev.md b/docs/webui-local-dev.md index a92e509..eacc739 100644 --- a/docs/webui-local-dev.md +++ b/docs/webui-local-dev.md @@ -48,7 +48,8 @@ that govern when a write path may open (#632, epic #631). | Path | Description | |------|-------------| | `/` | 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) | | `/api/queue` | JSON queue export with pagination metadata | | `/projects` | Project registry list (#427) | @@ -72,6 +73,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 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) Paste an LLM final report at `/audit` or POST JSON to `/api/audit`. The UI diff --git a/tests/test_webui_system_health.py b/tests/test_webui_system_health.py new file mode 100644 index 0000000..69cd1dd --- /dev/null +++ b/tests/test_webui_system_health.py @@ -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:secretpass@gitea.example.com/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:p@host.example.com/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:hunter2@gitea.example.com/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() diff --git a/webui/app.py b/webui/app.py index d0f832b..021b0db 100644 --- a/webui/app.py +++ b/webui/app.py @@ -29,6 +29,12 @@ from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as wo 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_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"}) _AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"}) @@ -62,16 +68,43 @@ async def home(_request: Request) -> HTMLResponse: 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 + started_at, uptime_seconds = process_uptime() return JSONResponse({ "status": "ok", "service": "mcp-control-plane-webui", "mode": "read-only-mvp", "timestamp": datetime.now(timezone.utc).isoformat(), "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: snapshot = load_queue_snapshot() return HTMLResponse(render_page(title="Queue", body_html=render_queue_page(snapshot))) @@ -263,6 +296,7 @@ def create_app(*, bind_host: str | None = None) -> Starlette: routes=[ Route("/", home, methods=["GET"]), Route("/health", health, methods=["GET"]), + Route(SYSTEM_HEALTH_API_PATH, api_system_health, methods=["GET"]), Route("/queue", queue, methods=["GET"]), Route("/api/queue", api_queue, methods=["GET"]), Route("/projects", projects, methods=["GET"]), diff --git a/webui/system_health.py b/webui/system_health.py new file mode 100644 index 0000000..a36438b --- /dev/null +++ b/webui/system_health.py @@ -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), + } From a6c15afec1ff3c154541bf65ecc196a317256c72 Mon Sep 17 00:00:00 2001 From: Jason Walker <913443@dadeschools.net> Date: Thu, 23 Jul 2026 02:38:19 -0400 Subject: [PATCH 2/4] fix: make cross-role allocations consumable by independent workers (Closes #843) Controller-created role=author allocations were owned by the allocating controller session with no authorized consume path for independent author workers. When the controller exited, the lease became stale_dead_process and required abandon/reassign instead of a usable handoff. - Mark cross-role apply with durable handoff provenance (pending) - Allow gitea_adopt_workflow_lease to consume pending handoffs by the required role without sharing controller session identity or requiring the controller process to remain alive - Atomically transfer assignment+lease ownership and set adopted_by_session_id with read-after-write evidence - Reject wrong-role, second, and terminal adoptions - Surface consume_allocation identifiers in process_work_queue results - Preserve same-role allocation and genuine abandon recovery behavior Co-Authored-By: Claude Opus 4.8 (1M context) --- allocator_service.py | 77 +++- control_plane_db.py | 172 +++++++- gitea_mcp_server.py | 9 +- lease_lifecycle.py | 187 ++++++++- tests/test_issue_843_cross_role_handoff.py | 449 +++++++++++++++++++++ 5 files changed, 872 insertions(+), 22 deletions(-) create mode 100644 tests/test_issue_843_cross_role_handoff.py diff --git a/allocator_service.py b/allocator_service.py index 84b5175..e6428f0 100644 --- a/allocator_service.py +++ b/allocator_service.py @@ -1115,6 +1115,12 @@ def allocate_next_work( "reasons": [ "dry-run only (apply=false); no assignment/lease created — " "call again with apply=true to reserve via control-plane DB" + + ( + "; after apply, the required-role worker consumes via " + "gitea_adopt_workflow_lease (#843)" + if mode == ALLOCATION_MODE_CROSS_ROLE and expected_role != role_norm + else "" + ) ], "skipped": [s.as_dict() for s in skipped], "terminal_pr": terminal_pr, @@ -1146,6 +1152,9 @@ def allocate_next_work( # Atomic reserve via #613 substrate. ttl = lease_ttl_seconds if lease_ttl_seconds is not None else None try: + cross_role_handoff = ( + mode == ALLOCATION_MODE_CROSS_ROLE and lease_role != role_norm + ) kwargs: dict[str, Any] = { "session_id": session_id, "role": lease_role, @@ -1157,7 +1166,8 @@ def allocate_next_work( "expected_head_sha": selected.head_sha, "allowed_actions": allowed, "forbidden_actions": forbidden, - "phase": "allocated", + # #843: mark cross-role allocations as awaiting independent consume + "phase": "awaiting_handoff" if cross_role_handoff else "allocated", } if ttl is not None: kwargs["lease_ttl_seconds"] = int(ttl) @@ -1237,7 +1247,52 @@ def allocate_next_work( "lease_role": lease_role, "source": "control_plane_db.assign_and_lease", } - return { + consume_allocation = None + if cross_role_handoff and result.lease_id: + # Durable handoff marker so independent required-role workers can + # consume without sharing the controller session (#843). + handoff_prov = { + "cross_role_handoff": True, + "handoff_status": "pending", + "allocating_session_id": session_id, + "allocating_role": role_norm, + "required_role": expected_role, + "required_profile": selection["required_profile"], + "required_namespace": selection["required_namespace"], + "assignment_id": result.assignment_id, + "lease_id": result.lease_id, + "allocation_mode": mode, + "adopted_by_session_id": None, + } + try: + db.attach_lease_provenance(result.lease_id, handoff_prov) + except ControlPlaneError: + # Still return assignment evidence; consume path may be unavailable + handoff_prov["attach_failed"] = True + consume_allocation = { + "tool": "gitea_adopt_workflow_lease", + "lease_id": result.lease_id, + "assignment_id": result.assignment_id, + "required_role": expected_role, + "required_profile": selection["required_profile"], + "required_namespace": selection["required_namespace"], + "handoff_status": "pending", + "controller_session_required": False, + "instructions": ( + f"From an independent {expected_role} session " + f"({selection['required_namespace']} / " + f"{selection['required_profile']}), call " + f"gitea_adopt_workflow_lease(lease_id={result.lease_id!r}) " + "to consume this controller allocation. The allocating " + "controller process does not need to remain alive. Wrong-role " + "and second-adoption attempts fail closed." + ), + } + lease_proof["cross_role_handoff"] = True + lease_proof["handoff_status"] = "pending" + lease_proof["consume_tool"] = "gitea_adopt_workflow_lease" + + out = { "success": True, "outcome": OUTCOME_ASSIGNED, "apply": True, @@ -1271,8 +1326,17 @@ def allocate_next_work( "lease_role": lease_role, "lease_proof": lease_proof, "selection_policy": SELECTION_POLICY, + "cross_role_handoff": bool(cross_role_handoff), }, - "next_valid_command": _next_command(lease_role, selected), + "next_valid_command": ( + ( + f"consume lease {result.lease_id} via gitea_adopt_workflow_lease " + f"as {expected_role}, then " + ) + + _next_command(lease_role, selected) + if cross_role_handoff + else _next_command(lease_role, selected) + ), "substrate": "control_plane_db", "file_lock_only": False, "comment_lease_only": False, @@ -1287,9 +1351,14 @@ def allocate_next_work( "downstream_note": ( "#612 incident bridge remains downstream of #600; " "allocator never assigns raw monitoring incidents; " - "controller routes only under cross_role (#840)" + "controller routes only under cross_role (#840); " + "cross-role assignments are consumable by independent " + "required-role workers via gitea_adopt_workflow_lease (#843)" ), } + if consume_allocation is not None: + out["consume_allocation"] = consume_allocation + return out def _next_command(role: str, c: WorkCandidate) -> str: diff --git a/control_plane_db.py b/control_plane_db.py index 8994a08..75af7c9 100644 --- a/control_plane_db.py +++ b/control_plane_db.py @@ -1637,11 +1637,13 @@ class ControlPlaneDB: provenance: dict[str, Any] | None = None, lease_ttl_seconds: int = DEFAULT_LEASE_TTL_SECONDS, ) -> dict[str, Any]: - """Transfer or refresh a lease with provenance (#601). + """Transfer or refresh a lease with provenance (#601 / #843). * Same owner + active → refresh (owner-resume). + * Cross-role handoff pending + matching required role → atomic consume + (even while the allocating controller session still "owns" the lease). * Expired/abandoned/released → create new assignment+lease with provenance. - * Active foreign → raise ForeignLeaseError (never silent steal). + * Active foreign (non-handoff) → raise ForeignLeaseError (never silent steal). """ now = _utc_now() now_s = _ts(now) @@ -1677,7 +1679,35 @@ class ControlPlaneDB: status = "expired" owner = lease["session_id"] - if status == "active" and owner != adopter_session_id: + # Parse durable provenance for cross-role handoff consume (#843). + lease_prov: dict[str, Any] = {} + if "provenance_json" in lease.keys() and lease["provenance_json"]: + try: + loaded = json.loads(lease["provenance_json"]) + if isinstance(loaded, dict): + lease_prov = loaded + except (TypeError, json.JSONDecodeError): + lease_prov = {} + handoff_pending = bool(lease_prov.get("cross_role_handoff")) and ( + str(lease_prov.get("handoff_status") or "pending").strip().lower() + == "pending" + ) + already_adopted = bool( + (lease["adopted_by_session_id"] if "adopted_by_session_id" in lease.keys() else None) + or lease_prov.get("adopted_by_session_id") + ) + required_role = str( + lease_prov.get("required_role") or lease["role"] or "" + ).strip().lower() + adopter_role = (role or "").strip().lower() + cross_role_consume = ( + handoff_pending + and not already_adopted + and status == "active" + and owner != adopter_session_id + ) + + if status == "active" and owner != adopter_session_id and not cross_role_consume: raise ForeignLeaseError( f"cannot adopt active foreign lease {lease_id} owned by {owner}" ) @@ -1761,6 +1791,142 @@ class ControlPlaneDB: "reasons": ["owner-resume: refreshed lease with provenance"], } + # #843: controller→required-role handoff consume (atomic, same lease_id) + if cross_role_consume: + if not required_role: + raise ControlPlaneError( + f"cross-role handoff lease {lease_id} missing required_role" + ) + if adopter_role != required_role: + raise ForeignLeaseError( + f"wrong role for cross-role handoff consume: " + f"required={required_role} adopter={adopter_role or 'none'} " + f"(fail closed)" + ) + # CAS: only transfer if still owned by allocating session and unadopted + cols = self._lease_columns(conn) + adopted_col_null = ( + "(adopted_by_session_id IS NULL OR adopted_by_session_id = '')" + if "adopted_by_session_id" in cols + else "1=1" + ) + cas = conn.execute( + f""" + UPDATE leases + SET session_id = ?, + heartbeat_at = ?, + expires_at = ?, + phase = ?, + role = ? + WHERE lease_id = ? + AND status = 'active' + AND session_id = ? + AND {adopted_col_null} + """, + ( + adopter_session_id, + now_s, + expires, + "adopted", + required_role, + lease_id, + owner, + ), + ) + if cas.rowcount != 1: + raise ForeignLeaseError( + f"cross-role handoff consume lost race for lease {lease_id} " + "(already adopted or no longer pending; fail closed)" + ) + if "adopted_from_session_id" in cols: + conn.execute( + """ + UPDATE leases + SET adopted_from_session_id = ?, adopted_by_session_id = ? + WHERE lease_id = ? + """, + (owner, adopter_session_id, lease_id), + ) + if "worktree_path" in cols and worktree_path: + conn.execute( + "UPDATE leases SET worktree_path = ? WHERE lease_id = ?", + (worktree_path, lease_id), + ) + if "owner_pid" in cols and owner_pid is not None: + conn.execute( + "UPDATE leases SET owner_pid = ? WHERE lease_id = ?", + (owner_pid, lease_id), + ) + if "expected_head_sha" in cols and expected_head_sha: + conn.execute( + "UPDATE leases SET expected_head_sha = ? WHERE lease_id = ?", + (expected_head_sha, lease_id), + ) + # Merge handoff provenance + caller provenance + merged = dict(lease_prov) + merged.update(provenance or {}) + merged["cross_role_handoff"] = True + merged["handoff_status"] = "adopted" + merged["adopted_from_session_id"] = owner + merged["adopted_by_session_id"] = adopter_session_id + merged["required_role"] = required_role + if "provenance_json" in cols: + conn.execute( + "UPDATE leases SET provenance_json = ? WHERE lease_id = ?", + (json.dumps(merged), lease_id), + ) + # Transfer active assignment ownership atomically + asn_cas = conn.execute( + """ + UPDATE assignments + SET session_id = ?, role = ? + WHERE lease_id = ? AND status = 'active' AND session_id = ? + """, + (adopter_session_id, required_role, lease_id, owner), + ) + if asn_cas.rowcount < 1: + # Fail closed: assignment must move with the lease + raise ControlPlaneError( + f"cross-role handoff: no active assignment for lease {lease_id} " + f"owned by {owner}" + ) + lease2 = conn.execute( + "SELECT * FROM leases WHERE lease_id = ?", (lease_id,) + ).fetchone() + asn = conn.execute( + """ + SELECT * FROM assignments + WHERE lease_id = ? AND status = 'active' + ORDER BY created_at DESC LIMIT 1 + """, + (lease_id,), + ).fetchone() + conn.execute( + """ + INSERT INTO events(work_item_id, event_type, message, created_at) + VALUES (?, 'lease_adopted', ?, ?) + """, + ( + lease["work_item_id"], + f"cross-role handoff: {adopter_session_id} consumed " + f"{lease_id} from {owner} as {required_role}", + now_s, + ), + ) + return { + "outcome": "adopted_cross_role_handoff", + "lease": dict(lease2) if lease2 else dict(lease), + "assignment": dict(asn) if asn else None, + "reasons": [ + "cross-role handoff: independent required-role worker consumed " + "controller allocation without abandonment" + ], + "adopted_by_session_id": adopter_session_id, + "adopted_from_session_id": owner, + "required_role": required_role, + "handoff_status": "adopted", + } + # Non-active: create new lease + assignment (transfer) new_lease_id = f"lease-{uuid.uuid4().hex[:16]}" new_asn_id = f"asn-{uuid.uuid4().hex[:16]}" diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index 94d6ec4..70835ab 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -21148,10 +21148,13 @@ def gitea_adopt_workflow_lease( remote: str = "dadeschools", host: str | None = None, ) -> dict: - """Adopt a control-plane lease through the sanctioned path (#601). + """Adopt a control-plane lease through the sanctioned path (#601 / #843). - Same-owner resume refreshes provenance. Foreign active leases are refused. - Expired leases may be reclaimed; provenance records adopted_from/by. + Same-owner resume refreshes provenance. Foreign active leases are refused + unless the lease is a pending controller cross-role handoff and the caller + holds the required role (independent consume without sharing the + controller session). Expired leases may be reclaimed; provenance records + adopted_from/by. Terminal (abandoned/released) leases cannot be adopted. """ read_block = _profile_operation_gate("gitea.read") if read_block: diff --git a/lease_lifecycle.py b/lease_lifecycle.py index 2049962..e081070 100644 --- a/lease_lifecycle.py +++ b/lease_lifecycle.py @@ -39,6 +39,7 @@ SAFE_RELEASE_OWNED = "release_owned" SAFE_STALE_PROMPT = "stale_prompt_lease" SAFE_UNKNOWN = "inspect_only" SAFE_NO_AUTHORITY = "file_or_comment_not_authoritative" +SAFE_CONSUME_CROSS_ROLE = "consume_cross_role_handoff" LEASE_STATUS_ACTIVE = "active" LEASE_STATUS_RELEASED = "released" @@ -250,6 +251,23 @@ def decide_safe_next_action( "same_owner": True, "also_allowed": [SAFE_ABANDON_ALLOWED, SAFE_RELEASE_OWNED], } + handoff = is_pending_cross_role_handoff({"lease": lease}) + if handoff: + return { + "safe_next_action": SAFE_CONSUME_CROSS_ROLE, + "reasons": [ + f"controller allocation pending handoff (freshness={status}); " + "required-role worker may consume without abandon/reassign; " + f"required_role={handoff['required_role']}" + ], + "block": False, + "same_owner": False, + "owner_session_id": owner, + "required_role": handoff["required_role"], + "cross_role_handoff": True, + "handoff_status": "pending", + "also_allowed": [SAFE_ABANDON_ALLOWED], + } return { "safe_next_action": SAFE_ABANDON_ALLOWED, "reasons": [ @@ -272,6 +290,24 @@ def decide_safe_next_action( } if not same_owner and status == "active": + # #843: pending cross-role handoff is consumable by required role + handoff = is_pending_cross_role_handoff({"lease": lease}) + if handoff: + return { + "safe_next_action": SAFE_CONSUME_CROSS_ROLE, + "reasons": [ + "controller cross-role allocation pending handoff; " + f"required_role={handoff['required_role']}; " + "consume via gitea_adopt_workflow_lease without " + "abandonment or sharing the controller session" + ], + "block": False, + "same_owner": False, + "owner_session_id": owner, + "required_role": handoff["required_role"], + "cross_role_handoff": True, + "handoff_status": "pending", + } return { "safe_next_action": SAFE_WAIT_FOREIGN, "reasons": [ @@ -440,6 +476,84 @@ def list_active_leases( } + +def parse_lease_provenance(lease_or_state: Mapping[str, Any] | None) -> dict[str, Any]: + """Return durable lease provenance dict (empty when absent/unparseable).""" + if not lease_or_state: + return {} + if "provenance" in lease_or_state and isinstance(lease_or_state.get("provenance"), dict): + return dict(lease_or_state["provenance"]) + raw = None + if "provenance_json" in lease_or_state: + raw = lease_or_state.get("provenance_json") + elif "lease" in lease_or_state and isinstance(lease_or_state.get("lease"), Mapping): + raw = lease_or_state["lease"].get("provenance_json") + if not raw: + return {} + if isinstance(raw, dict): + return dict(raw) + try: + loaded = json.loads(raw) + except (TypeError, json.JSONDecodeError): + return {} + return dict(loaded) if isinstance(loaded, dict) else {} + + +def is_pending_cross_role_handoff( + state: Mapping[str, Any] | None, +) -> dict[str, Any] | None: + """Return handoff evidence when a controller allocation awaits consume (#843). + + A pending handoff is identified by durable provenance written at + cross-role apply time — not by title heuristics or session-id guessing. + """ + if not state: + return None + lease = state.get("lease") if isinstance(state.get("lease"), Mapping) else state + if not isinstance(lease, Mapping): + return None + status = str(lease.get("status") or "").strip().lower() + if status in (LEASE_STATUS_ABANDONED, LEASE_STATUS_RELEASED, LEASE_STATUS_EXPIRED): + return None + prov = parse_lease_provenance(state) + if not prov and isinstance(lease, Mapping): + prov = parse_lease_provenance(lease) + if not prov.get("cross_role_handoff"): + return None + handoff_status = str(prov.get("handoff_status") or "pending").strip().lower() + if handoff_status != "pending": + return None + adopted_by = ( + lease.get("adopted_by_session_id") + or prov.get("adopted_by_session_id") + or "" + ) + if str(adopted_by).strip(): + return None + required_role = str( + prov.get("required_role") or lease.get("role") or "" + ).strip().lower() + if not required_role: + return None + return { + "cross_role_handoff": True, + "handoff_status": "pending", + "required_role": required_role, + "allocating_session_id": str( + prov.get("allocating_session_id") or lease.get("session_id") or "" + ), + "allocating_role": str(prov.get("allocating_role") or "controller"), + "lease_id": str(lease.get("lease_id") or ""), + "assignment_id": ( + str(state["assignment"]["assignment_id"]) + if isinstance(state.get("assignment"), Mapping) + and state["assignment"].get("assignment_id") + else None + ), + "provenance": prov, + } + + def adopt_lease( db: cpd.ControlPlaneDB, *, @@ -463,11 +577,8 @@ def adopt_lease( owner = str(lease.get("session_id") or "") same_owner = owner == str(adopter_session_id) - if freshness["freshness"] == "active" and not same_owner: - raise LeaseLifecycleError( - f"refusing to steal active foreign lease {lease_id} owned by " - f"{owner} (fail closed)" - ) + handoff = is_pending_cross_role_handoff(state) + adopter_role = (role or "").strip().lower() if freshness["freshness"] in ("abandoned", "released"): raise LeaseLifecycleError( @@ -475,13 +586,40 @@ def adopt_lease( "(fail closed)" ) - # Expired or stale: require abandon-style safety before ownership transfer - # when not same owner; same owner may reclaim. - if not same_owner and freshness["freshness"] in ( + if handoff and not same_owner: + # Terminal statuses already rejected above. Freshness may be + # active OR stale_dead_process (controller exited) — both are + # consumable without abandonment when handoff is still pending. + if freshness["freshness"] not in ( + "active", + "stale_dead_process", + "stale_missing_worktree", + ): + raise LeaseLifecycleError( + f"lease {lease_id} freshness={freshness['freshness']}; " + "terminal or non-active allocation cannot be handoff-consumed " + "(fail closed)" + ) + required = handoff["required_role"] + if adopter_role != required: + raise LeaseLifecycleError( + f"wrong role for cross-role handoff consume of {lease_id}: " + f"required={required} adopter={adopter_role or 'none'} " + "(fail closed)" + ) + reason = "cross-role-handoff-consume" + elif freshness["freshness"] == "active" and not same_owner: + raise LeaseLifecycleError( + f"refusing to steal active foreign lease {lease_id} owned by " + f"{owner} (fail closed)" + ) + elif not same_owner and freshness["freshness"] in ( "expired", "stale_dead_process", "stale_missing_worktree", ): + # Expired or stale (non-handoff): require abandon-style safety before + # ownership transfer when not same owner; same owner may reclaim. if not operator_authorized and freshness["freshness"] == "expired": # Deterministic reclaim of expired foreign lease is allowed # without operator flag (sanctioned expire reclaim). @@ -492,6 +630,9 @@ def adopt_lease( f"lease {lease_id} freshness={freshness['freshness']}; " "use abandon with proof before foreign adopt (fail closed)" ) + reason = "sanctioned-reclaim-adopt" + else: + reason = "owner-resume-adopt" if same_owner else "sanctioned-reclaim-adopt" provenance = build_adopt_provenance( adopted_from_session_id=owner, @@ -504,10 +645,14 @@ def adopt_lease( worktree_path=worktree_path, expected_head_sha=expected_head_sha or lease.get("expected_head_sha"), prior_lease_id=lease_id, - reason=( - "owner-resume-adopt" if same_owner else "sanctioned-reclaim-adopt" - ), + reason=reason, ) + if handoff and not same_owner: + provenance["cross_role_handoff"] = True + provenance["handoff_status"] = "adopted" + provenance["required_role"] = handoff["required_role"] + provenance["allocating_session_id"] = handoff["allocating_session_id"] + provenance["allocating_role"] = handoff["allocating_role"] result = db.adopt_lease( lease_id=lease_id, @@ -518,7 +663,7 @@ def adopt_lease( owner_pid=owner_pid if owner_pid is not None else os.getpid(), provenance=provenance, ) - return { + out = { "success": True, "outcome": result.get("outcome"), "same_owner": same_owner, @@ -531,6 +676,24 @@ def adopt_lease( "comment_lease_only": False, "reasons": result.get("reasons") or [], } + if handoff and not same_owner: + out["cross_role_handoff"] = True + out["handoff_status"] = "adopted" + out["required_role"] = handoff["required_role"] + out["adopted_by_session_id"] = adopter_session_id + out["adopted_from_session_id"] = owner + lease_row = result.get("lease") or {} + if isinstance(lease_row, Mapping): + out["read_after_write"] = { + "lease_id": lease_row.get("lease_id"), + "session_id": lease_row.get("session_id"), + "role": lease_row.get("role"), + "status": lease_row.get("status"), + "adopted_by_session_id": lease_row.get("adopted_by_session_id"), + "adopted_from_session_id": lease_row.get("adopted_from_session_id"), + "phase": lease_row.get("phase"), + } + return out def release_lease( diff --git a/tests/test_issue_843_cross_role_handoff.py b/tests/test_issue_843_cross_role_handoff.py new file mode 100644 index 0000000..d616f13 --- /dev/null +++ b/tests/test_issue_843_cross_role_handoff.py @@ -0,0 +1,449 @@ +"""Cross-role allocation handoff consumable by independent workers (#843). + +Regression coverage for the controller→required-role consume path: + +* controller allocates author work; independent author adopts successfully +* author adoption succeeds after allocating controller process exits +* author adoption without sharing controller session identity +* wrong-role adoption rejected +* concurrent/second adoption rejected without state corruption +* terminal allocation adoption rejected +* successful adoption produces authoritative ownership evidence +* genuine abandoned-lease recovery remains valid +* process_work_queue / allocate results include consume identifiers +* same-role allocation behavior remains compatible +""" + +from __future__ import annotations + +import os +import tempfile +import unittest +from datetime import timedelta +from unittest.mock import patch + +from allocator_service import ( + ALLOCATION_MODE_CROSS_ROLE, + ALLOCATION_MODE_ROLE_SCOPED, + OUTCOME_ASSIGNED, + ROLE_AUTHOR, + ROLE_CONTROLLER, + ROLE_REVIEWER, + WorkCandidate, + allocate_next_work, +) +from control_plane_db import ControlPlaneDB, ForeignLeaseError, _ts, _utc_now +import lease_lifecycle as ll + + +class CrossRoleHandoffTest(unittest.TestCase): + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory() + self.db_path = os.path.join(self._tmp.name, "cp.sqlite3") + self.db = ControlPlaneDB(self.db_path) + self.db.upsert_session( + session_id="ctrl-session", + role="controller", + profile="prgs-controller", + pid=99999999, # dead-looking pid + ) + self.db.upsert_session( + session_id="author-worker", + role="author", + profile="prgs-author", + pid=os.getpid(), + ) + self.db.upsert_session( + session_id="author-worker-2", + role="author", + profile="prgs-author", + pid=os.getpid(), + ) + self.db.upsert_session( + session_id="reviewer-worker", + role="reviewer", + profile="prgs-reviewer", + pid=os.getpid(), + ) + self.wt = self._tmp.name + + def tearDown(self) -> None: + self._tmp.cleanup() + + def _ready_issue(self, number: int = 843, title: str = "handoff target") -> WorkCandidate: + return WorkCandidate( + kind="issue", + number=number, + labels=("status:ready", "type:bug"), + title=title, + priority=20, + ) + + def _controller_allocate(self, number: int = 843, **kwargs): + defaults = dict( + db=self.db, + session_id="ctrl-session", + role=ROLE_CONTROLLER, + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + candidates=[self._ready_issue(number)], + apply=True, + profile_name="prgs-controller", + username="controller-user", + allocation_mode=ALLOCATION_MODE_CROSS_ROLE, + ) + defaults.update(kwargs) + return allocate_next_work(**defaults) + + def test_controller_allocates_author_independent_author_adopts(self) -> None: + res = self._controller_allocate() + self.assertEqual(res["outcome"], OUTCOME_ASSIGNED) + self.assertEqual(res["required_role"], ROLE_AUTHOR) + self.assertIn("consume_allocation", res) + consume = res["consume_allocation"] + self.assertEqual(consume["tool"], "gitea_adopt_workflow_lease") + self.assertEqual(consume["required_role"], ROLE_AUTHOR) + self.assertFalse(consume["controller_session_required"]) + lid = res["assignment"]["lease_id"] + self.assertEqual(consume["lease_id"], lid) + self.assertIn(lid, res["next_valid_command"]) + + adopted = ll.adopt_lease( + self.db, + lease_id=lid, + adopter_session_id="author-worker", + role=ROLE_AUTHOR, + worktree_path=self.wt, + ) + self.assertTrue(adopted["success"]) + self.assertEqual(adopted["outcome"], "adopted_cross_role_handoff") + self.assertEqual(adopted["adopted_by_session_id"], "author-worker") + self.assertEqual(adopted["adopted_from_session_id"], "ctrl-session") + raw = adopted["read_after_write"] + self.assertEqual(raw["session_id"], "author-worker") + self.assertEqual(raw["adopted_by_session_id"], "author-worker") + self.assertEqual(raw["status"], "active") + self.assertEqual(raw["phase"], "adopted") + + # Authoritative re-read + state = self.db.get_lease_workflow_state(lid) + self.assertEqual(state["lease"]["session_id"], "author-worker") + self.assertEqual(state["lease"]["adopted_by_session_id"], "author-worker") + self.assertEqual(state["assignment"]["session_id"], "author-worker") + self.assertEqual(state["provenance"]["handoff_status"], "adopted") + + def test_author_adoption_after_controller_process_exits(self) -> None: + res = self._controller_allocate(number=900) + lid = res["assignment"]["lease_id"] + # Force owner_pid dead + freshness stale_dead_process + import sqlite3 + + conn = sqlite3.connect(self.db_path) + try: + conn.execute( + "UPDATE leases SET owner_pid = 99999999 WHERE lease_id = ?", + (lid,), + ) + conn.commit() + finally: + conn.close() + state = self.db.get_lease_workflow_state(lid) + fr = ll.classify_lease_freshness( + state["lease"], pid_checker=lambda _p: False + ) + self.assertEqual(fr["freshness"], "stale_dead_process") + + adopted = ll.adopt_lease( + self.db, + lease_id=lid, + adopter_session_id="author-worker", + role=ROLE_AUTHOR, + worktree_path=self.wt, + ) + self.assertEqual(adopted["outcome"], "adopted_cross_role_handoff") + self.assertEqual(adopted["adopted_by_session_id"], "author-worker") + # No abandon required + state2 = self.db.get_lease_workflow_state(lid) + self.assertEqual(state2["lease"]["status"], "active") + self.assertNotEqual(state2["lease"]["status"], "abandoned") + + def test_adoption_without_sharing_controller_session_identity(self) -> None: + res = self._controller_allocate(number=901) + lid = res["assignment"]["lease_id"] + adopted = ll.adopt_lease( + self.db, + lease_id=lid, + adopter_session_id="author-worker", + role=ROLE_AUTHOR, + worktree_path=self.wt, + ) + self.assertNotEqual(adopted["adopted_by_session_id"], "ctrl-session") + self.assertFalse(adopted["same_owner"]) + self.assertEqual(adopted["adopted_from_session_id"], "ctrl-session") + + def test_wrong_role_adoption_rejected(self) -> None: + res = self._controller_allocate(number=902) + lid = res["assignment"]["lease_id"] + with self.assertRaises(ll.LeaseLifecycleError) as ctx: + ll.adopt_lease( + self.db, + lease_id=lid, + adopter_session_id="reviewer-worker", + role=ROLE_REVIEWER, + ) + self.assertIn("wrong role", str(ctx.exception).lower()) + # State unchanged + state = self.db.get_lease_workflow_state(lid) + self.assertEqual(state["lease"]["session_id"], "ctrl-session") + self.assertIsNone(state["lease"].get("adopted_by_session_id") or None) + self.assertEqual(state["provenance"]["handoff_status"], "pending") + + def test_second_adoption_rejected_without_corruption(self) -> None: + res = self._controller_allocate(number=903) + lid = res["assignment"]["lease_id"] + first = ll.adopt_lease( + self.db, + lease_id=lid, + adopter_session_id="author-worker", + role=ROLE_AUTHOR, + worktree_path=self.wt, + ) + self.assertEqual(first["outcome"], "adopted_cross_role_handoff") + with self.assertRaises(ll.LeaseLifecycleError): + ll.adopt_lease( + self.db, + lease_id=lid, + adopter_session_id="author-worker-2", + role=ROLE_AUTHOR, + worktree_path=self.wt, + ) + state = self.db.get_lease_workflow_state(lid) + self.assertEqual(state["lease"]["session_id"], "author-worker") + self.assertEqual(state["lease"]["adopted_by_session_id"], "author-worker") + self.assertEqual(state["assignment"]["session_id"], "author-worker") + self.assertEqual(state["lease"]["status"], "active") + + def test_terminal_allocation_adoption_rejected(self) -> None: + res = self._controller_allocate(number=904) + lid = res["assignment"]["lease_id"] + # Abandon as terminal + proof = ll.AbandonProof( + dead_process=True, + missing_worktree=True, + no_open_pr=True, + no_live_mutation_risk=True, + owner_pid=99999999, + worktree_path="/nonexistent/for-843", + ) + # Attach dead pid / missing wt for abandon eligibility + import sqlite3 + + conn = sqlite3.connect(self.db_path) + try: + conn.execute( + "UPDATE leases SET owner_pid = 99999999, worktree_path = ? WHERE lease_id = ?", + ("/nonexistent/for-843", lid), + ) + conn.commit() + finally: + conn.close() + abandoned = ll.abandon_lease( + self.db, + lease_id=lid, + requester_session_id="author-worker", + proof=proof, + ) + self.assertEqual(abandoned["outcome"], "abandoned") + with self.assertRaises(ll.LeaseLifecycleError) as ctx: + ll.adopt_lease( + self.db, + lease_id=lid, + adopter_session_id="author-worker", + role=ROLE_AUTHOR, + ) + self.assertIn("abandoned", str(ctx.exception).lower()) + + def test_successful_adoption_read_after_write_ownership(self) -> None: + res = self._controller_allocate(number=905) + lid = res["assignment"]["lease_id"] + adopted = ll.adopt_lease( + self.db, + lease_id=lid, + adopter_session_id="author-worker", + role=ROLE_AUTHOR, + worktree_path=self.wt, + ) + raw = adopted["read_after_write"] + self.assertEqual(raw["lease_id"], lid) + self.assertEqual(raw["session_id"], "author-worker") + self.assertEqual(raw["adopted_by_session_id"], "author-worker") + self.assertEqual(raw["adopted_from_session_id"], "ctrl-session") + # Re-fetch proves durable write + state = self.db.get_lease_workflow_state(lid) + self.assertEqual(state["lease"]["session_id"], raw["session_id"]) + self.assertEqual( + state["lease"]["adopted_by_session_id"], raw["adopted_by_session_id"] + ) + + def test_genuine_abandoned_recovery_still_valid(self) -> None: + """Same-role author lease abandoned remains reclaimable via abandon path.""" + same = allocate_next_work( + self.db, + session_id="author-worker", + role=ROLE_AUTHOR, + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + candidates=[self._ready_issue(906, "same-role")], + apply=True, + profile_name="prgs-author", + username="author-user", + allocation_mode=ALLOCATION_MODE_ROLE_SCOPED, + ) + self.assertEqual(same["outcome"], OUTCOME_ASSIGNED) + lid = same["assignment"]["lease_id"] + import sqlite3 + + conn = sqlite3.connect(self.db_path) + try: + conn.execute( + "UPDATE leases SET owner_pid = 99999999, worktree_path = ? WHERE lease_id = ?", + ("/nonexistent/same-role", lid), + ) + conn.commit() + finally: + conn.close() + proof = ll.AbandonProof( + dead_process=True, + missing_worktree=True, + no_open_pr=True, + no_live_mutation_risk=True, + owner_pid=99999999, + worktree_path="/nonexistent/same-role", + ) + abandoned = ll.abandon_lease( + self.db, + lease_id=lid, + requester_session_id="author-worker-2", + proof=proof, + ) + self.assertEqual(abandoned["outcome"], "abandoned") + # Foreign author cannot handoff-consume an abandoned non-handoff lease + with self.assertRaises(ll.LeaseLifecycleError): + ll.adopt_lease( + self.db, + lease_id=lid, + adopter_session_id="author-worker-2", + role=ROLE_AUTHOR, + ) + # Reclaim path still works for expired/abandoned after force-expire + reclaimed = ll.reclaim_expired_lease( + self.db, + lease_id=lid, + session_id="author-worker-2", + role=ROLE_AUTHOR, + worktree_path=self.wt, + ) + self.assertEqual(reclaimed["outcome"], "reclaimed") + self.assertEqual(reclaimed["assignment"]["session_id"], "author-worker-2") + + def test_allocate_payload_includes_consume_identifiers(self) -> None: + res = self._controller_allocate(number=907) + self.assertIn("consume_allocation", res) + c = res["consume_allocation"] + for key in ( + "tool", + "lease_id", + "assignment_id", + "required_role", + "required_profile", + "required_namespace", + "instructions", + "handoff_status", + ): + self.assertIn(key, c) + self.assertEqual(c["required_namespace"], "gitea-author") + self.assertEqual(c["required_profile"], "prgs-author") + self.assertIn("gitea_adopt_workflow_lease", c["instructions"]) + self.assertTrue(res["lease_proof"]["cross_role_handoff"]) + self.assertEqual(res["lease_proof"]["handoff_status"], "pending") + + def test_same_role_allocation_remains_compatible(self) -> None: + res = allocate_next_work( + self.db, + session_id="author-worker", + role=ROLE_AUTHOR, + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + candidates=[self._ready_issue(908)], + apply=True, + profile_name="prgs-author", + username="author-user", + ) + self.assertEqual(res["outcome"], OUTCOME_ASSIGNED) + self.assertNotIn("consume_allocation", res) + lid = res["assignment"]["lease_id"] + state = self.db.get_lease_workflow_state(lid) + # No cross-role handoff provenance + prov = state.get("provenance") or {} + self.assertFalse(prov.get("cross_role_handoff")) + # Owner resume still works + resume = ll.adopt_lease( + self.db, + lease_id=lid, + adopter_session_id="author-worker", + role=ROLE_AUTHOR, + worktree_path=self.wt, + ) + self.assertTrue(resume["same_owner"]) + self.assertEqual(resume["outcome"], "adopted_owner_resume") + + def test_inspect_points_required_role_at_consume(self) -> None: + res = self._controller_allocate(number=909) + lid = res["assignment"]["lease_id"] + decision = ll.inspect_lease( + self.db, lid, caller_session_id="author-worker" + ) + self.assertEqual( + decision["safe_next_action"], ll.SAFE_CONSUME_CROSS_ROLE + ) + self.assertFalse(decision["block"]) + self.assertEqual(decision["required_role"], ROLE_AUTHOR) + + def test_db_cas_rejects_concurrent_second_consume(self) -> None: + res = self._controller_allocate(number=910) + lid = res["assignment"]["lease_id"] + # First consume via DB layer directly + first = self.db.adopt_lease( + lease_id=lid, + adopter_session_id="author-worker", + role=ROLE_AUTHOR, + worktree_path=self.wt, + provenance={ + "cross_role_handoff": True, + "handoff_status": "adopted", + "required_role": "author", + }, + ) + self.assertEqual(first["outcome"], "adopted_cross_role_handoff") + # Second CAS must fail + with self.assertRaises(ForeignLeaseError): + self.db.adopt_lease( + lease_id=lid, + adopter_session_id="author-worker-2", + role=ROLE_AUTHOR, + worktree_path=self.wt, + provenance={ + "cross_role_handoff": True, + "handoff_status": "pending", + "required_role": "author", + }, + ) + state = self.db.get_lease_workflow_state(lid) + self.assertEqual(state["lease"]["session_id"], "author-worker") + + +if __name__ == "__main__": + unittest.main() From 5eb89f883074cf8ab56461767188e9454ad04a98 Mon Sep 17 00:00:00 2001 From: Jason Walker <913443@dadeschools.net> Date: Thu, 23 Jul 2026 03:41:28 -0400 Subject: [PATCH 3/4] fix: bind cross-role handoff consume role to authenticated profile (Closes #843) Review #515 F1: gitea_adopt_workflow_lease trusted a caller-supplied role ((role or active_role)), so any namespace holding gitea.read could consume an author-only cross-role handoff by passing role="author". - Derive the adopter role authoritatively from the active profile; reject any supplied role that does not exactly match (no silent accept). - Pass the profile-derived role and authoritative profile/namespace context to lease_lifecycle.adopt_lease; validate handoff provenance required_profile/required_namespace against it (fail closed). - Fail closed when the profile role cannot be derived (no author default). - Add MCP-boundary regression tests: reviewer/merger profiles cannot consume an author handoff via role="author"; the legitimate author profile still consumes; foreign required_profile rejected. --- gitea_mcp_server.py | 55 ++++- lease_lifecycle.py | 35 +++- tests/test_issue_843_cross_role_handoff.py | 233 +++++++++++++++++++++ 3 files changed, 318 insertions(+), 5 deletions(-) diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index 70835ab..fd5db07 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -21155,6 +21155,14 @@ def gitea_adopt_workflow_lease( holds the required role (independent consume without sharing the controller session). Expired leases may be reclaimed; provenance records adopted_from/by. Terminal (abandoned/released) leases cannot be adopted. + + #843 F1: the adopter role is derived authoritatively from the active + authenticated profile — never from caller input. A supplied ``role`` that + does not exactly match the profile-derived role is rejected (no silent + accept or reinterpretation), and handoff provenance ``required_profile`` / + ``required_namespace`` restrictions are validated against the same + authoritative caller context. Caller-supplied role/profile/namespace can + never grant authority. """ read_block = _profile_operation_gate("gitea.read") if read_block: @@ -21163,25 +21171,64 @@ def gitea_adopt_workflow_lease( "reasons": read_block, "permission_report": _permission_block_report("gitea.read"), } + profile = get_profile() + profile_name = (profile.get("profile_name") or "").strip() or "session" + active_role = (_profile_role_kind(profile) or "").strip().lower() + if not active_role: + return { + "success": False, + "outcome": "blocked", + "mutation_performed": False, + "reasons": [ + "active profile role could not be derived authoritatively; " + "refusing lease adoption (fail closed, #843)" + ], + "lease_id": lease_id, + "authoritative_source": "control_plane_db", + "file_lock_only": False, + "comment_lease_only": False, + } + if role is not None and str(role).strip(): + supplied_role = str(role).strip().lower() + if supplied_role != active_role: + return { + "success": False, + "outcome": "blocked", + "mutation_performed": False, + "profile_role_kind": active_role, + "supplied_role": supplied_role, + "reasons": [ + f"caller-supplied role '{supplied_role}' does not match " + f"the authenticated profile-derived role '{active_role}'; " + "caller-supplied role/profile/namespace can never grant " + "authority (fail closed, #843)" + ], + "lease_id": lease_id, + "authoritative_source": "control_plane_db", + "file_lock_only": False, + "comment_lease_only": False, + } db, errs = _control_plane_db_or_error() if db is None: return {"success": False, "reasons": errs} - profile = get_profile() - profile_name = (profile.get("profile_name") or "").strip() or "session" - active_role = _profile_role_kind(profile) or "author" sid = (session_id or "").strip() or ( f"{profile_name}-{os.getpid()}-{uuid.uuid4().hex[:8]}" ) + adopter_namespace = allocator_service.DEFAULT_ROLE_NAMESPACES.get( + active_role, f"gitea-{active_role}" + ) try: return lease_lifecycle.adopt_lease( db, lease_id=lease_id, adopter_session_id=sid, - role=(role or active_role).strip() or "author", + role=active_role, worktree_path=worktree_path, expected_head_sha=expected_head_sha, owner_pid=os.getpid(), operator_authorized=bool(operator_authorized), + adopter_profile_name=profile_name, + adopter_namespace=adopter_namespace, ) except (lease_lifecycle.LeaseLifecycleError, control_plane_db.ControlPlaneError) as exc: return { diff --git a/lease_lifecycle.py b/lease_lifecycle.py index e081070..c42be14 100644 --- a/lease_lifecycle.py +++ b/lease_lifecycle.py @@ -564,8 +564,17 @@ def adopt_lease( expected_head_sha: str | None = None, owner_pid: int | None = None, operator_authorized: bool = False, + adopter_profile_name: str | None = None, + adopter_namespace: str | None = None, ) -> dict[str, Any]: - """Sanctioned adopt path with provenance; never silent foreign steal.""" + """Sanctioned adopt path with provenance; never silent foreign steal. + + #843 F1: for a pending cross-role handoff, ``role`` must be the + authoritative profile-derived role supplied by the MCP boundary — never + caller-asserted authority. When the handoff provenance declares + ``required_profile`` / ``required_namespace`` and the caller context is + provided, both are validated exactly; a mismatch fails closed. + """ state = db.get_lease_workflow_state(lease_id) if not state: raise LeaseLifecycleError( @@ -607,6 +616,30 @@ def adopt_lease( f"required={required} adopter={adopter_role or 'none'} " "(fail closed)" ) + # #843 F1: provenance profile/namespace restrictions are validated + # against the authoritative caller context when declared. Caller + # input can never widen authority; a mismatch fails closed. + handoff_prov = handoff.get("provenance") or {} + required_profile = str( + handoff_prov.get("required_profile") or "" + ).strip() + if required_profile and adopter_profile_name is not None: + if str(adopter_profile_name).strip() != required_profile: + raise LeaseLifecycleError( + f"wrong profile for cross-role handoff consume of " + f"{lease_id}: required_profile={required_profile} " + f"adopter_profile={adopter_profile_name} (fail closed)" + ) + required_namespace = str( + handoff_prov.get("required_namespace") or "" + ).strip() + if required_namespace and adopter_namespace is not None: + if str(adopter_namespace).strip() != required_namespace: + raise LeaseLifecycleError( + f"wrong namespace for cross-role handoff consume of " + f"{lease_id}: required_namespace={required_namespace} " + f"adopter_namespace={adopter_namespace} (fail closed)" + ) reason = "cross-role-handoff-consume" elif freshness["freshness"] == "active" and not same_owner: raise LeaseLifecycleError( diff --git a/tests/test_issue_843_cross_role_handoff.py b/tests/test_issue_843_cross_role_handoff.py index d616f13..4fe0028 100644 --- a/tests/test_issue_843_cross_role_handoff.py +++ b/tests/test_issue_843_cross_role_handoff.py @@ -445,5 +445,238 @@ class CrossRoleHandoffTest(unittest.TestCase): self.assertEqual(state["lease"]["session_id"], "author-worker") +class MCPBoundaryAdoptRoleBindingTest(unittest.TestCase): + """#843 F1: MCP-boundary role binding for ``gitea_adopt_workflow_lease``. + + The library-level wrong-role test calls ``lease_lifecycle.adopt_lease`` + directly. These tests prove the MCP entry point derives the adopter role + authoritatively from the active authenticated profile and rejects any + caller-supplied role that disagrees, so a reviewer/merger profile cannot + consume an author handoff by passing ``role="author"``. + """ + + AUTHOR_PROFILE = { + "profile_name": "prgs-author", + "role": "author", + "allowed_operations": [ + "gitea.read", + "gitea.pr.create", + "gitea.branch.push", + ], + "forbidden_operations": [], + } + REVIEWER_PROFILE = { + "profile_name": "prgs-reviewer", + "role": "reviewer", + "allowed_operations": [ + "gitea.read", + "gitea.pr.review", + "gitea.pr.approve", + "gitea.pr.request_changes", + ], + "forbidden_operations": ["gitea.pr.create", "gitea.branch.push"], + } + MERGER_PROFILE = { + "profile_name": "prgs-merger", + "role": "merger", + "allowed_operations": ["gitea.read", "gitea.pr.merge"], + "forbidden_operations": ["gitea.pr.create", "gitea.branch.push"], + } + FOREIGN_AUTHOR_PROFILE = { + "profile_name": "dadeschools-author", + "role": "author", + "allowed_operations": [ + "gitea.read", + "gitea.pr.create", + "gitea.branch.push", + ], + "forbidden_operations": [], + } + + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory() + self.db_path = os.path.join(self._tmp.name, "cp.sqlite3") + self.db = ControlPlaneDB(self.db_path) + self.db.upsert_session( + session_id="ctrl-session", + role="controller", + profile="prgs-controller", + pid=99999999, + ) + self.db.upsert_session( + session_id="author-worker", + role="author", + profile="prgs-author", + pid=os.getpid(), + ) + self.wt = self._tmp.name + + def tearDown(self) -> None: + self._tmp.cleanup() + + def _ready_issue(self, number: int) -> WorkCandidate: + return WorkCandidate( + kind="issue", + number=number, + labels=("status:ready", "type:bug"), + title="handoff target", + priority=20, + ) + + def _handoff_lease(self, number: int = 843) -> str: + res = allocate_next_work( + db=self.db, + session_id="ctrl-session", + role=ROLE_CONTROLLER, + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + candidates=[self._ready_issue(number)], + apply=True, + profile_name="prgs-controller", + username="controller-user", + allocation_mode=ALLOCATION_MODE_CROSS_ROLE, + ) + self.assertEqual(res["outcome"], OUTCOME_ASSIGNED) + self.assertEqual(res["required_role"], ROLE_AUTHOR) + return res["assignment"]["lease_id"] + + def _call_adopt_tool(self, profile: dict, **kwargs): + import gitea_mcp_server as mcp_server + + with ( + patch.object(mcp_server, "get_profile", return_value=profile), + patch.object( + mcp_server, + "_control_plane_db_or_error", + return_value=(self.db, []), + ), + ): + return mcp_server.gitea_adopt_workflow_lease( + remote="prgs", **kwargs + ) + + def _assert_handoff_untouched(self, lease_id: str) -> None: + state = self.db.get_lease_workflow_state(lease_id) + self.assertEqual(state["lease"]["session_id"], "ctrl-session") + self.assertIsNone(state["lease"].get("adopted_by_session_id") or None) + self.assertEqual(state["lease"]["status"], "active") + self.assertEqual(state["provenance"]["handoff_status"], "pending") + + def test_reviewer_profile_cannot_consume_author_handoff_via_role_author( + self, + ) -> None: + lid = self._handoff_lease(920) + result = self._call_adopt_tool( + self.REVIEWER_PROFILE, + lease_id=lid, + session_id="reviewer-worker", + role="author", + worktree_path=self.wt, + ) + self.assertFalse(result["success"]) + self.assertEqual(result["outcome"], "blocked") + self.assertEqual(result["profile_role_kind"], "reviewer") + self.assertEqual(result["supplied_role"], "author") + self.assertIn("does not match", result["reasons"][0]) + self._assert_handoff_untouched(lid) + + def test_merger_profile_cannot_consume_author_handoff_via_role_author( + self, + ) -> None: + lid = self._handoff_lease(921) + result = self._call_adopt_tool( + self.MERGER_PROFILE, + lease_id=lid, + session_id="merger-worker", + role="author", + worktree_path=self.wt, + ) + self.assertFalse(result["success"]) + self.assertEqual(result["outcome"], "blocked") + self.assertEqual(result["profile_role_kind"], "merger") + self._assert_handoff_untouched(lid) + + def test_reviewer_profile_rejected_without_role_argument(self) -> None: + """Even without a spoofed role, the profile-derived role binds.""" + lid = self._handoff_lease(922) + result = self._call_adopt_tool( + self.REVIEWER_PROFILE, + lease_id=lid, + session_id="reviewer-worker", + worktree_path=self.wt, + ) + self.assertFalse(result["success"]) + self.assertEqual(result["outcome"], "blocked") + self.assertIn("wrong role", result["reasons"][0].lower()) + self._assert_handoff_untouched(lid) + + def test_author_profile_mismatching_supplied_role_rejected(self) -> None: + lid = self._handoff_lease(923) + result = self._call_adopt_tool( + self.AUTHOR_PROFILE, + lease_id=lid, + session_id="author-worker", + role="reviewer", + worktree_path=self.wt, + ) + self.assertFalse(result["success"]) + self.assertEqual(result["outcome"], "blocked") + self.assertEqual(result["profile_role_kind"], "author") + self.assertEqual(result["supplied_role"], "reviewer") + self._assert_handoff_untouched(lid) + + def test_foreign_profile_name_rejected_for_author_handoff(self) -> None: + """Provenance required_profile binds even when the role matches.""" + lid = self._handoff_lease(924) + result = self._call_adopt_tool( + self.FOREIGN_AUTHOR_PROFILE, + lease_id=lid, + session_id="foreign-author-worker", + worktree_path=self.wt, + ) + self.assertFalse(result["success"]) + self.assertEqual(result["outcome"], "blocked") + self.assertIn("wrong profile", result["reasons"][0].lower()) + self._assert_handoff_untouched(lid) + + def test_author_profile_consumes_author_handoff(self) -> None: + lid = self._handoff_lease(925) + result = self._call_adopt_tool( + self.AUTHOR_PROFILE, + lease_id=lid, + session_id="author-worker", + role="author", + worktree_path=self.wt, + ) + self.assertTrue(result["success"]) + self.assertEqual(result["outcome"], "adopted_cross_role_handoff") + self.assertEqual(result["adopted_by_session_id"], "author-worker") + self.assertEqual(result["adopted_from_session_id"], "ctrl-session") + state = self.db.get_lease_workflow_state(lid) + self.assertEqual(state["lease"]["session_id"], "author-worker") + self.assertEqual( + state["lease"]["adopted_by_session_id"], "author-worker" + ) + self.assertEqual(state["assignment"]["session_id"], "author-worker") + self.assertEqual(state["provenance"]["handoff_status"], "adopted") + + def test_author_profile_consumes_author_handoff_without_role_argument( + self, + ) -> None: + lid = self._handoff_lease(926) + result = self._call_adopt_tool( + self.AUTHOR_PROFILE, + lease_id=lid, + session_id="author-worker", + worktree_path=self.wt, + ) + self.assertTrue(result["success"]) + self.assertEqual(result["outcome"], "adopted_cross_role_handoff") + state = self.db.get_lease_workflow_state(lid) + self.assertEqual(state["lease"]["session_id"], "author-worker") + self.assertEqual(state["lease"]["role"], "author") + + if __name__ == "__main__": unittest.main() From c3f282ba44a788f24fc281ce39edddbfe0fdea8e Mon Sep 17 00:00:00 2001 From: Jason Walker <913443@dadeschools.net> Date: Thu, 23 Jul 2026 04:56:58 -0400 Subject: [PATCH 4/4] fix: exclude epic and child-only containers from allocator selection (Closes #844) Epics and parent issues whose body delegates implementation to children are skipped before ranking with structured reason epic_or_child_only_container. Title-only "epic" mentions without body/label evidence remain eligible. Co-Authored-By: Claude Opus 4.8 (1M context) --- allocator_service.py | 104 +++++++- gitea_mcp_server.py | 1 + ...test_allocator_epic_container_exclusion.py | 243 ++++++++++++++++++ 3 files changed, 347 insertions(+), 1 deletion(-) create mode 100644 tests/test_allocator_epic_container_exclusion.py diff --git a/allocator_service.py b/allocator_service.py index 84b5175..50f4bde 100644 --- a/allocator_service.py +++ b/allocator_service.py @@ -53,6 +53,8 @@ OUTCOME_CANDIDATE_SET_DRIFT = "candidate_set_drift" SKIP_CLAIMED_BY_OTHER_SESSION = "claimed_by_other_session" # #776: controller-supplied pre-rank exclusion. SKIP_EXCLUDED_BY_CONTROLLER = "excluded_by_controller" +# #844: epic / child-only implementation container (pre-rank). +SKIP_EPIC_OR_CHILD_ONLY_CONTAINER = "epic_or_child_only_container" # Ownership verdicts for a live claim on a candidate (#765). OWNERSHIP_OWN = "own" @@ -130,6 +132,39 @@ ROLE_ACTIONS: dict[str, tuple[tuple[str, ...], tuple[str, ...]]] = { } +# Body phrases that prove an issue is an implementation container, not a +# unit of direct author work (#844). Matched case-insensitively against the +# issue body. Title alone is never sufficient (ordinary issues may mention +# "epic" incidentally). +_CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = ( + "implementation is delivered via child issues only", + "implementation is delivered through child issues only", + "implementation is delivered via child issues", + "implementation is delivered through child issues", + "do not implement product features in this epic", + "do not implement product features in this epic issue itself", + "no product feature implementation is claimed complete solely on this epic", + "implementable child issues remain independently eligible", + "owns the product roadmap and linkage", + "this epic owns the product roadmap", + "coordination container", + "child-only container", + "implementation is delegated to child", +) + +# Explicit epic / umbrella labels (structured evidence preferred over title). +_EPIC_LABELS: frozenset[str] = frozenset( + { + "type:epic", + "epic", + "kind:epic", + "scope:epic", + "type:umbrella", + "umbrella", + } +) + + @dataclass class WorkCandidate: """One assignable Gitea issue or PR presented to the allocator.""" @@ -139,6 +174,7 @@ class WorkCandidate: state: str = "open" labels: tuple[str, ...] = () title: str = "" + body: str = "" priority: int = 0 head_sha: str | None = None # Routing signals (callers derive from Gitea / review feedback). @@ -158,6 +194,7 @@ class WorkCandidate: self.labels = tuple( str(x).strip().lower() for x in (self.labels or ()) if str(x).strip() ) + self.body = str(self.body or "") if self.kind not in WORK_KINDS: raise InvalidWorkKindError( f"candidate kind '{self.kind}' is not assignable; only " @@ -171,6 +208,7 @@ class WorkCandidate: "state": self.state, "labels": list(self.labels), "title": self.title, + "body": self.body, "priority": self.priority, "head_sha": self.head_sha, "request_changes_current_head": self.request_changes_current_head, @@ -184,6 +222,51 @@ class WorkCandidate: } +def classify_epic_or_child_only_container( + c: WorkCandidate, +) -> tuple[bool, str | None]: + """Return whether *c* is an epic / child-only implementation container (#844). + + Exclusion uses structured evidence first (labels, body scope language). + A bare title containing the word "epic" is **not** enough — ordinary + implementable issues may mention epics incidentally. A title that is + explicitly prefixed ``Epic:`` only counts when the body also proves + child-only / no-direct-implementation scope (or an epic label is present). + + PRs are never classified as containers here (they already have a head). + """ + if c.kind != "issue": + return False, None + + labels = set(c.labels) + epic_label = sorted(labels & _EPIC_LABELS) + body_l = (c.body or "").lower() + title = (c.title or "").strip() + title_l = title.lower() + + body_hits = [m for m in _CHILD_ONLY_BODY_MARKERS if m in body_l] + title_epic_prefix = title_l.startswith("epic:") or title_l.startswith("epic ") + + if epic_label: + detail = f"label={epic_label[0]}" + if body_hits: + detail = f"{detail}; body_marker={body_hits[0]!r}" + return True, detail + + if body_hits: + # Body proves child-only / umbrella scope. Title "Epic:" is corroborating + # but not required — containers without the word still exclude. + detail = f"body_marker={body_hits[0]!r}" + if title_epic_prefix: + detail = f"title_epic_prefix; {detail}" + return True, detail + + # Title-only "Epic:" without body scope evidence is insufficient (#844 AC: + # eligibility does not rely solely on the word "Epic" in a title). + # Similarly, incidental "epic" mid-title without markers stays eligible. + return False, None + + @dataclass class SkipRecord: kind: str @@ -850,7 +933,8 @@ def allocate_next_work( ownership_defects: list[dict[str, Any]] = [] controller_excluded: list[dict[str, Any]] = [] - # #776 AC2: remove excluded numbers *before* ranking / selection / lease. + # #776 AC2 + #844: remove excluded numbers *and* epic/child-only containers + # *before* ranking / selection / lease so they never receive assignments. rankable: list[WorkCandidate] = [] for c in candidates: if int(c.number) in exclude_set: @@ -929,6 +1013,23 @@ def allocate_next_work( }, } continue + # #844: epics / child-only containers are never direct implement targets. + is_container, container_detail = classify_epic_or_child_only_container(c) + if is_container: + detail = container_detail or "epic or child-only container" + reason = ( + f"{c.kind}#{c.number} {SKIP_EPIC_OR_CHILD_ONLY_CONTAINER}: " + f"{detail}; implementation is delegated to child issues" + ) + skipped.append( + SkipRecord( + c.kind, + c.number, + reason, + SKIP_EPIC_OR_CHILD_ONLY_CONTAINER, + ) + ) + continue rankable.append(c) ordered = sort_candidates(rankable) @@ -1341,6 +1442,7 @@ def candidate_from_dict(data: dict[str, Any]) -> WorkCandidate: state=str(data.get("state") or "open"), labels=tuple(data.get("labels") or ()), title=str(data.get("title") or ""), + body=str(data.get("body") or ""), priority=priority, head_sha=data.get("head_sha"), request_changes_current_head=bool(data.get("request_changes_current_head")), diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index 94d6ec4..d63b38d 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -19912,6 +19912,7 @@ def _allocator_candidates_from_gitea( state="open", labels=tuple(labels), title=title, + body=body, priority=20 if "status:ready" in labels else 1, blocked=blocked, dependency_unmet=dep_unmet, diff --git a/tests/test_allocator_epic_container_exclusion.py b/tests/test_allocator_epic_container_exclusion.py new file mode 100644 index 0000000..0b9fa78 --- /dev/null +++ b/tests/test_allocator_epic_container_exclusion.py @@ -0,0 +1,243 @@ +"""Allocator epic / child-only container pre-rank exclusion (#844). + +Covers: +* Issue #631-shaped child-only epic is excluded before ranking. +* Implementable child issues remain eligible and can be selected. +* Ordinary issues that merely mention "epic" in title/body are not excluded. +* Excluded containers never receive assignments or workflow leases. +* Structured skip reason ``epic_or_child_only_container`` is reported. +""" + +from __future__ import annotations + +import os +import tempfile +import unittest + +from allocator_service import ( + OUTCOME_ASSIGNED, + OUTCOME_PREVIEW, + SKIP_EPIC_OR_CHILD_ONLY_CONTAINER, + WorkCandidate, + allocate_next_work, + classify_epic_or_child_only_container, +) +from control_plane_db import ControlPlaneDB + +REMOTE = "prgs" +ORG = "Scaled-Tech-Consulting" +REPO = "Gitea-Tools" + +# Minimal body mirroring issue #631 authoritative scope language. +_EPIC_631_BODY = """ +## Scope (umbrella) + +This epic owns the **product roadmap and linkage** for the Web Console. +Implementation is delivered via child issues only. + +## Explicit non-goals + +* Do not implement product features in this epic issue itself. +* No product feature implementation is claimed complete solely on this epic. +""" + +_CHILD_BODY = """ +## Problem + +Operators need a workflow-event timeline model for Phase 1. + +## Acceptance criteria + +- [ ] Timeline model API exists +""" + + +def _issue( + number: int, + *, + title: str = "", + body: str = "", + labels: tuple[str, ...] = ("status:ready", "type:feature"), + priority: int = 20, +) -> WorkCandidate: + return WorkCandidate( + kind="issue", + number=number, + state="open", + labels=labels, + title=title or f"issue {number}", + body=body, + priority=priority, + ) + + +class ClassifyEpicContainerTest(unittest.TestCase): + def test_631_shaped_body_and_title_is_container(self) -> None: + c = _issue( + 631, + title="Epic: MCP Control Plane Web Console", + body=_EPIC_631_BODY, + ) + is_c, detail = classify_epic_or_child_only_container(c) + self.assertTrue(is_c) + self.assertIsNotNone(detail) + self.assertIn("body_marker", detail or "") + + def test_body_markers_without_epic_title(self) -> None: + c = _issue( + 900, + title="Control plane roadmap tracker", + body="Implementation is delivered via child issues only.", + ) + is_c, _ = classify_epic_or_child_only_container(c) + self.assertTrue(is_c) + + def test_epic_label_alone_is_container(self) -> None: + c = _issue( + 901, + title="Roadmap linkage", + body="Track children.", + labels=("status:ready", "type:epic"), + ) + is_c, detail = classify_epic_or_child_only_container(c) + self.assertTrue(is_c) + self.assertIn("type:epic", detail or "") + + def test_title_epic_prefix_alone_not_container(self) -> None: + """Title-only 'Epic:' without body scope evidence stays eligible (#844).""" + c = _issue( + 902, + title="Epic: something mentioned only in title", + body="Implement a concrete fix for the allocator skip list.", + ) + is_c, detail = classify_epic_or_child_only_container(c) + self.assertFalse(is_c) + self.assertIsNone(detail) + + def test_incidental_epic_word_not_container(self) -> None: + c = _issue( + 903, + title="Document epic handoff conventions", + body=( + "Update the docs so implementable issues that mention an epic " + "remain independently executable." + ), + ) + is_c, _ = classify_epic_or_child_only_container(c) + self.assertFalse(is_c) + + def test_prs_never_classified(self) -> None: + pr = WorkCandidate( + kind="pr", + number=10, + state="open", + title="Epic: fake", + body="Implementation is delivered via child issues only.", + head_sha="a" * 40, + priority=5, + ) + is_c, _ = classify_epic_or_child_only_container(pr) + self.assertFalse(is_c) + + +class AllocateEpicContainerExclusionTest(unittest.TestCase): + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory() + self.addCleanup(self._tmp.cleanup) + self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3")) + + def _alloc(self, candidates, **kwargs): + defaults = dict( + session_id="sess-844", + role="author", + remote=REMOTE, + org=ORG, + repo=REPO, + profile_name="prgs-author", + username="jcwalker3", + claims={}, + apply=False, + ) + defaults.update(kwargs) + return allocate_next_work(self.db, candidates=candidates, **defaults) + + def test_631_shaped_epic_excluded_child_selected(self) -> None: + epic = _issue( + 631, + title="Epic: MCP Control Plane Web Console", + body=_EPIC_631_BODY, + ) + child = _issue( + 637, + title="Web Console: Workflow-event timeline model (Phase 1)", + body=_CHILD_BODY, + ) + res = self._alloc([epic, child], apply=False) + self.assertTrue(res["success"], res) + self.assertEqual(res["outcome"], OUTCOME_PREVIEW) + self.assertEqual(res["selected"]["number"], 637) + skipped = {s["number"]: s for s in res["skipped"]} + self.assertIn(631, skipped) + self.assertEqual( + skipped[631]["reason_code"], SKIP_EPIC_OR_CHILD_ONLY_CONTAINER + ) + self.assertIn(SKIP_EPIC_OR_CHILD_ONLY_CONTAINER, skipped[631]["reason"]) + + def test_container_cannot_receive_assignment_or_lease(self) -> None: + epic = _issue( + 631, + title="Epic: MCP Control Plane Web Console", + body=_EPIC_631_BODY, + ) + res = self._alloc([epic], apply=True) + self.assertTrue(res["success"], res) + # Only container present → no safe work; never assigned_work. + self.assertNotEqual(res["outcome"], OUTCOME_ASSIGNED) + self.assertIsNone(res.get("assignment")) + self.assertIsNone(res.get("selected")) + skipped = {s["number"]: s for s in res["skipped"]} + self.assertEqual( + skipped[631]["reason_code"], SKIP_EPIC_OR_CHILD_ONLY_CONTAINER + ) + # No lease row for the epic. + leases = self.db.list_active_leases( + remote=REMOTE, org=ORG, repo=REPO + ) if hasattr(self.db, "list_active_leases") else [] + # Prefer generic inventory if available. + if not leases and hasattr(self.db, "list_leases"): + leases = self.db.list_leases(remote=REMOTE, org=ORG, repo=REPO) + for lease in leases or []: + work_number = lease.get("work_number") if isinstance(lease, dict) else None + self.assertNotEqual(work_number, 631) + + def test_incidental_epic_title_remains_eligible(self) -> None: + ordinary = _issue( + 700, + title="Document epic handoff conventions", + body="Write runbook text about epic vs child issues.", + ) + res = self._alloc([ordinary], apply=False) + self.assertTrue(res["success"], res) + self.assertEqual(res["selected"]["number"], 700) + self.assertEqual(res["skipped"], []) + + def test_apply_selects_child_not_epic(self) -> None: + epic = _issue( + 631, + title="Epic: MCP Control Plane Web Console", + body=_EPIC_631_BODY, + ) + child = _issue( + 637, + title="Web Console: Workflow-event timeline model (Phase 1)", + body=_CHILD_BODY, + ) + res = self._alloc([epic, child], apply=True) + self.assertTrue(res["success"], res) + self.assertEqual(res["outcome"], OUTCOME_ASSIGNED) + self.assertEqual(res["selected"]["number"], 637) + self.assertEqual(res["assignment"]["work_number"], 637) + + +if __name__ == "__main__": + unittest.main()