From 5494696227b84148b374412e9b38d58c1eccfca5 Mon Sep 17 00:00:00 2001 From: jcwalker3 Date: Wed, 22 Jul 2026 15:59:30 -0500 Subject: [PATCH 1/2] 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 648d9464ba088cb850d037847355d67874c57466 Mon Sep 17 00:00:00 2001 From: Jason Walker <913443@dadeschools.net> Date: Thu, 23 Jul 2026 00:05:35 -0400 Subject: [PATCH 2/2] fix: authoritative cross-role generic queue allocation (Closes #840) Add controller-owned cross_role allocation mode that inspects the full queue and returns one selection with required role/profile/action and lease evidence. Document process_work_queue routing, normalize controller role metadata, and keep the dashboard explanatory only. Co-Authored-By: Claude Opus 4.8 (1M context) --- allocator_service.py | 334 +++++++++++-- gitea_mcp_server.py | 47 +- namespace_workspace_binding.py | 24 +- role_session_router.py | 35 ++ task_capability_map.py | 19 +- tests/test_cross_role_queue_allocation.py | 581 ++++++++++++++++++++++ workflow_dashboard.py | 8 +- 7 files changed, 992 insertions(+), 56 deletions(-) create mode 100644 tests/test_cross_role_queue_allocation.py diff --git a/allocator_service.py b/allocator_service.py index 40787ce..84b5175 100644 --- a/allocator_service.py +++ b/allocator_service.py @@ -78,6 +78,33 @@ VALID_ROLES = frozenset( {ROLE_AUTHOR, ROLE_REVIEWER, ROLE_MERGER, ROLE_RECONCILER, ROLE_CONTROLLER} ) +# Allocation modes (#840). +# role_scoped: only candidates whose expected role matches the caller role. +# cross_role: controller-owned generic queue selection — inspect full queue, +# rank/eligibility canonically, return one selection naming the required +# downstream role/profile. Controller routes; it does not perform mutations. +ALLOCATION_MODE_ROLE_SCOPED = "role_scoped" +ALLOCATION_MODE_CROSS_ROLE = "cross_role" +VALID_ALLOCATION_MODES = frozenset( + {ALLOCATION_MODE_ROLE_SCOPED, ALLOCATION_MODE_CROSS_ROLE} +) + +# Default execution-profile / MCP-namespace names for each role. +DEFAULT_ROLE_PROFILES: dict[str, str] = { + ROLE_AUTHOR: "prgs-author", + ROLE_REVIEWER: "prgs-reviewer", + ROLE_MERGER: "prgs-merger", + ROLE_RECONCILER: "prgs-reconciler", + ROLE_CONTROLLER: "prgs-controller", +} +DEFAULT_ROLE_NAMESPACES: dict[str, str] = { + ROLE_AUTHOR: "gitea-author", + ROLE_REVIEWER: "gitea-reviewer", + ROLE_MERGER: "gitea-merger", + ROLE_RECONCILER: "gitea-reconciler", + ROLE_CONTROLLER: "gitea-controller", +} + # Default action matrices by role (mutation gate will re-check). ROLE_ACTIONS: dict[str, tuple[tuple[str, ...], tuple[str, ...]]] = { ROLE_AUTHOR: ( @@ -259,6 +286,126 @@ def normalize_role(role: str | None, *, profile_name: str | None = None) -> str: ) +def resolve_allocation_mode( + role: str, + allocation_mode: str | None = None, +) -> str: + """Resolve allocation mode; controller defaults to cross_role (#840).""" + raw = (allocation_mode or "").strip().lower() + if raw: + if raw not in VALID_ALLOCATION_MODES: + raise ControlPlaneError( + f"unknown allocation_mode {allocation_mode!r}; expected one of " + f"{sorted(VALID_ALLOCATION_MODES)}" + ) + return raw + if role == ROLE_CONTROLLER: + return ALLOCATION_MODE_CROSS_ROLE + return ALLOCATION_MODE_ROLE_SCOPED + + +def required_profile_for_role( + role: str, + *, + profile_name: str | None = None, +) -> str: + """Map a required role to the canonical execution profile name.""" + role_norm = (role or "").strip().lower() + # Preserve remote/env prefix from the active profile when present + # (e.g. dadeschools-author → dadeschools-reviewer). + active = (profile_name or "").strip() + if active: + lower = active.lower() + for token in ("author", "reviewer", "merger", "reconciler", "controller"): + if lower.endswith(f"-{token}") or lower == token: + prefix = active[: -len(token)].rstrip("-") + if prefix: + return f"{prefix}-{role_norm}" + return role_norm + return DEFAULT_ROLE_PROFILES.get(role_norm, f"prgs-{role_norm}") + + +def required_namespace_for_role( + role: str, + *, + profile_name: str | None = None, +) -> str: + """Map a required role to the canonical MCP namespace name.""" + role_norm = (role or "").strip().lower() + profile = required_profile_for_role(role_norm, profile_name=profile_name) + # Namespace is typically gitea-; keep stable mapping when profile is + # non-prgs (still gitea- for isolation). + return DEFAULT_ROLE_NAMESPACES.get(role_norm, f"gitea-{role_norm}") + + +def selected_action_for_candidate(c: WorkCandidate, required_role: str) -> str: + """Canonical next action for the selected work under *required_role*.""" + role = (required_role or "").strip().lower() + if role == ROLE_AUTHOR: + if c.kind == "pr" and c.request_changes_current_head: + return "address_pr_change_requests" + if c.kind == "pr": + return "update_pr" + return "implement" + if role == ROLE_REVIEWER: + if c.approval_stale: + return "re_review" + return "review" + if role == ROLE_MERGER: + return "merge" + if role == ROLE_RECONCILER: + if c.approval_contaminated: + return "reconcile_contaminated_approval" + return "reconcile" + if role == ROLE_CONTROLLER: + return "diagnose" + return "process" + + +def build_selection_dict( + selected: WorkCandidate, + *, + active_role: str, + required_role: str, + profile_name: str | None = None, + allocation_mode: str, +) -> dict[str, Any]: + """Authoritative single selection payload for allocator results (#840).""" + action = selected_action_for_candidate(selected, required_role) + req_profile = required_profile_for_role( + required_role, profile_name=profile_name + ) + req_ns = required_namespace_for_role( + required_role, profile_name=profile_name + ) + return { + "kind": selected.kind, + "number": selected.number, + "title": selected.title, + "labels": list(selected.labels), + "head_sha": selected.head_sha, + "priority": selected.priority, + "expected_role_next": required_role, + "required_role": required_role, + "selected_action": action, + "action": action, + "required_profile": req_profile, + "required_namespace": req_ns, + "pinned": { + "kind": selected.kind, + "number": selected.number, + "head_sha": selected.head_sha, + "issue_number": selected.number if selected.kind == "issue" else None, + "pr_number": selected.number if selected.kind == "pr" else None, + }, + "reason_selected": ( + f"highest-priority eligible candidate under allocation_mode=" + f"'{allocation_mode}' (active_role={active_role}, " + f"required_role={required_role}, action={action})" + ), + } + + def expected_role_for_candidate(c: WorkCandidate) -> str: """ADR §5.3 routing: which role should take this work next.""" if c.kind == "pr": @@ -289,6 +436,7 @@ def classify_skip( role: str, terminal_pr: int | None, claim_ownership: str | None = None, + allocation_mode: str | None = None, ) -> str | None: """Return skip reason, or None if candidate is selectable for *role*. @@ -297,7 +445,12 @@ def classify_skip( and unknown claims are excluded so one session's in-progress task can never blockade the queue for a different controller; ``own`` stays selectable so a controller can resume its own work. + + *allocation_mode* (#840): ``cross_role`` (controller default) ranks the full + queue and selects the highest-priority eligible item for any downstream + role. ``role_scoped`` retains prior role-match filtering. """ + mode = resolve_allocation_mode(role, allocation_mode) if c.state in ("merged", "closed"): return f"{c.kind}#{c.number} is {c.state}; never assign" if c.blocked or "status:blocked" in c.labels: @@ -322,34 +475,58 @@ def classify_skip( if c.kind == "pr" and not (c.head_sha or "").strip(): return f"pr#{c.number} missing head_sha pin" + expected = expected_role_for_candidate(c) + # Terminal path first: when an active terminal PR exists, only that PR - # (or controller diagnosis) is assignable for review-path roles. + # is assignable for review-path roles (or for work whose expected role is + # review/merge under cross_role selection). if terminal_pr is not None and c.kind == "pr" and c.number != terminal_pr: - if role in (ROLE_REVIEWER, ROLE_MERGER): + terminal_roles = (ROLE_REVIEWER, ROLE_MERGER) + if mode == ALLOCATION_MODE_CROSS_ROLE: + if expected in terminal_roles: + return ( + f"pr#{c.number} skipped: active terminal-review lock on " + f"PR #{terminal_pr} must be resolved first" + ) + elif role in terminal_roles: return ( f"pr#{c.number} skipped: active terminal-review lock on " f"PR #{terminal_pr} must be resolved first" ) - expected = expected_role_for_candidate(c) - if role == ROLE_CONTROLLER: - # Controller may inspect anything but only assigns diagnosis targets - # when contaminated / blocked. + if mode == ALLOCATION_MODE_CROSS_ROLE: + # Cross-role controller selection: eligibility only — no active-role + # match filter. The selection payload names required_role. + pass + elif role == ROLE_CONTROLLER: + # Legacy diagnosis-only controller path (role_scoped): only reconciler- + # needed targets. Prefer cross_role for generic queue allocation. if expected == ROLE_RECONCILER or c.blocked: return None - return f"{c.kind}#{c.number} does not require controller (expected {expected})" - - if role != expected: + return ( + f"{c.kind}#{c.number} does not require controller " + f"(expected {expected})" + ) + elif role != expected: return ( f"{c.kind}#{c.number} expects role '{expected}', active role is '{role}'" ) # Ready-gate for issues: prefer status:ready when labels present. + # Applies for author-bound work in both modes (cross_role only gates + # author-expected issues so reconciler/reviewer PRs stay selectable). if c.kind == "issue" and c.labels: - if "status:ready" not in c.labels and "status:in-progress" not in c.labels: - # Allow unlabeled open issues; only skip explicit non-ready states. - if any(l.startswith("status:") for l in c.labels): - return f"issue#{c.number} not status:ready ({','.join(c.labels)})" + gate_role = expected if mode == ALLOCATION_MODE_CROSS_ROLE else role + if gate_role in (ROLE_AUTHOR, ROLE_CONTROLLER): + if ( + "status:ready" not in c.labels + and "status:in-progress" not in c.labels + ): + if any(l.startswith("status:") for l in c.labels): + return ( + f"issue#{c.number} not status:ready " + f"({','.join(c.labels)})" + ) return None @@ -501,12 +678,19 @@ def allocate_next_work( claims: Mapping[tuple[str, int], dict[str, Any]] | None = None, exclude_issue_numbers: Sequence[int] | None = None, expected_candidate_set_fingerprint: str | None = None, + allocation_mode: str | None = None, ) -> dict[str, Any]: """Select and optionally reserve the next work unit via control-plane DB. *apply=False* (default): dry-run selection only — no lease/assignment. *apply=True*: atomic ``assign_and_lease`` for the selected candidate. + *allocation_mode* (#840): ``cross_role`` (default for controller) inspects + the complete queue and returns one authoritative selection naming the + required downstream role/profile/action. ``role_scoped`` keeps prior + per-role filtering. Controller routes only — never grants author/reviewer/ + merger/reconciler mutation rights to the controller session. + *exclude_issue_numbers* (#776): numbers removed before ranking. Omitted / empty preserves prior behavior. @@ -541,6 +725,19 @@ def allocate_next_work( "substrate": "control_plane_db", } + try: + mode = resolve_allocation_mode(role_norm, allocation_mode) + except ControlPlaneError as exc: + return { + "success": False, + "outcome": OUTCOME_ROLE_INELIGIBLE, + "reasons": [str(exc)], + "skipped": [], + "assignment": None, + "substrate": "control_plane_db", + "allocation_mode": (allocation_mode or "").strip() or None, + } + session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}" try: db.upsert_session( @@ -748,6 +945,7 @@ def allocate_next_work( role=role_norm, terminal_pr=terminal_pr, claim_ownership=ownership, + allocation_mode=mode, ) if reason: is_claim_skip = SKIP_CLAIMED_BY_OTHER_SESSION in reason @@ -828,6 +1026,10 @@ def allocate_next_work( "outcome": outcome, "apply": bool(apply), "role": role_norm, + "allocation_mode": mode, + "routing_role": role_norm, + "required_role": None, + "selected_action": None, "profile_name": profile_name, "username": username, "session_id": session_id, @@ -840,6 +1042,12 @@ def allocate_next_work( "skipped": [s.as_dict() for s in skipped], "terminal_pr": terminal_pr, "assignment": None, + "allocation_evidence": { + "mode": "empty", + "allocation_mode": mode, + "lease_created": False, + "selection_policy": SELECTION_POLICY, + }, "substrate": "control_plane_db", "file_lock_only": False, "comment_lease_only": False, @@ -852,25 +1060,37 @@ def allocate_next_work( "owner_session_id": owner_session_id, "downstream_note": ( "#612 incident bridge remains downstream of #600; " - "allocator never assigns raw monitoring incidents" + "allocator never assigns raw monitoring incidents; " + "controller routes only under cross_role (#840)" ), } expected_role = expected_role_for_candidate(selected) - allowed, forbidden = role_actions(role_norm) - selection = { - "kind": selected.kind, - "number": selected.number, - "title": selected.title, - "labels": list(selected.labels), - "head_sha": selected.head_sha, - "priority": selected.priority, - "expected_role_next": expected_role, - "reason_selected": ( - f"highest-priority candidate for role '{role_norm}' " - f"(expected_role={expected_role})" - ), - } + # Cross-role: lease/action matrix follows the required downstream role so + # evidence names the worker that must act. Controller session still owns + # the routing decision; mutation isolation is enforced by role gates on + # mutation tools (controller profile lacks author/review/merge ops). + lease_role = ( + expected_role if mode == ALLOCATION_MODE_CROSS_ROLE else role_norm + ) + allowed, forbidden = role_actions(lease_role) + # Controller must never receive mutation-class rights via cross-role apply. + if role_norm == ROLE_CONTROLLER: + ctrl_allowed, ctrl_forbidden = role_actions(ROLE_CONTROLLER) + # Keep controller session capability evidence separate from lease_role. + controller_allowed_actions = ctrl_allowed + controller_forbidden_actions = ctrl_forbidden + else: + controller_allowed_actions = allowed + controller_forbidden_actions = forbidden + + selection = build_selection_dict( + selected, + active_role=role_norm, + required_role=expected_role, + profile_name=profile_name, + allocation_mode=mode, + ) if not apply: return { @@ -878,6 +1098,12 @@ def allocate_next_work( "outcome": OUTCOME_PREVIEW, "apply": False, "role": role_norm, + "allocation_mode": mode, + "routing_role": role_norm, + "required_role": expected_role, + "selected_action": selection["selected_action"], + "required_profile": selection["required_profile"], + "required_namespace": selection["required_namespace"], "profile_name": profile_name, "username": username, "session_id": session_id, @@ -893,6 +1119,12 @@ def allocate_next_work( "skipped": [s.as_dict() for s in skipped], "terminal_pr": terminal_pr, "assignment": None, + "allocation_evidence": { + "mode": "preview", + "allocation_mode": mode, + "lease_created": False, + "selection_policy": SELECTION_POLICY, + }, "substrate": "control_plane_db", "file_lock_only": False, "comment_lease_only": False, @@ -902,9 +1134,12 @@ def allocate_next_work( "controller_excluded": list(controller_excluded), "exclude_issue_numbers": list(exclude_nums), "candidate_set_fingerprint": cas_fp, + "controller_allowed_actions": list(controller_allowed_actions), + "controller_forbidden_actions": list(controller_forbidden_actions), "downstream_note": ( "#612 incident bridge remains downstream of #600; " - "allocator never assigns raw monitoring incidents" + "allocator never assigns raw monitoring incidents; " + "controller routes only under cross_role (#840)" ), } @@ -913,7 +1148,7 @@ def allocate_next_work( try: kwargs: dict[str, Any] = { "session_id": session_id, - "role": role_norm, + "role": lease_role, "remote": remote, "org": org, "repo": repo, @@ -992,11 +1227,27 @@ def allocate_next_work( } # assigned + lease_proof = { + "assignment_id": result.assignment_id, + "lease_id": result.lease_id, + "expires_at": result.expires_at, + "expected_head_sha": result.expected_head_sha, + "allowed_actions": list(result.allowed_actions), + "forbidden_actions": list(result.forbidden_actions), + "lease_role": lease_role, + "source": "control_plane_db.assign_and_lease", + } return { "success": True, "outcome": OUTCOME_ASSIGNED, "apply": True, "role": role_norm, + "allocation_mode": mode, + "routing_role": role_norm, + "required_role": expected_role, + "selected_action": selection["selected_action"], + "required_profile": selection["required_profile"], + "required_namespace": selection["required_namespace"], "profile_name": profile_name, "username": username, "session_id": session_id, @@ -1012,16 +1263,16 @@ def allocate_next_work( "skipped": [s.as_dict() for s in skipped], "terminal_pr": terminal_pr, "assignment": result.as_dict(), - "lease_proof": { - "assignment_id": result.assignment_id, - "lease_id": result.lease_id, - "expires_at": result.expires_at, - "expected_head_sha": result.expected_head_sha, - "allowed_actions": list(result.allowed_actions), - "forbidden_actions": list(result.forbidden_actions), - "source": "control_plane_db.assign_and_lease", + "lease_proof": lease_proof, + "allocation_evidence": { + "mode": "assigned", + "allocation_mode": mode, + "lease_created": True, + "lease_role": lease_role, + "lease_proof": lease_proof, + "selection_policy": SELECTION_POLICY, }, - "next_valid_command": _next_command(role_norm, selected), + "next_valid_command": _next_command(lease_role, selected), "substrate": "control_plane_db", "file_lock_only": False, "comment_lease_only": False, @@ -1031,9 +1282,12 @@ def allocate_next_work( "controller_excluded": list(controller_excluded), "exclude_issue_numbers": list(exclude_nums), "candidate_set_fingerprint": cas_fp, + "controller_allowed_actions": list(controller_allowed_actions), + "controller_forbidden_actions": list(controller_forbidden_actions), "downstream_note": ( "#612 incident bridge remains downstream of #600; " - "allocator never assigns raw monitoring incidents" + "allocator never assigns raw monitoring incidents; " + "controller routes only under cross_role (#840)" ), } diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index acd56ca..94d6ec4 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -234,12 +234,25 @@ def _effective_workspace_role() -> str: def _profile_role_kind(profile: dict) -> str: - """Resolve a profile's declared role before inferring from permissions.""" - role = (profile.get("role") or profile.get("role_kind") or "").strip() + """Resolve a profile's declared role before inferring from permissions. + + Declared ``role`` / ``role_kind`` always wins so a controller profile is + never reclassified as reconciler from permission inference (#840). + """ + role = (profile.get("role") or profile.get("role_kind") or "").strip().lower() if role: + # Normalize aliases / case. + if "control" in role: + return "controller" return role profile_name = (profile.get("profile_name") or "").strip().lower() - for candidate in ("reconciler", "merger", "reviewer", "author"): + for candidate in ( + "controller", + "reconciler", + "merger", + "reviewer", + "author", + ): if candidate in profile_name: return candidate return _role_kind( @@ -15538,7 +15551,8 @@ def mcp_get_control_plane_guide( profile = get_profile() allowed = profile["allowed_operations"] forbidden = profile["forbidden_operations"] - role = _role_kind(allowed, forbidden) + # Prefer declared profile role so controller is not mislabeled reconciler (#840). + role = _profile_role_kind(profile) username = _authenticated_username(h) identity = { @@ -15597,6 +15611,16 @@ def mcp_get_control_plane_guide( "user, and merging requires explicit operator authorization plus the " "'MERGE PR ' confirmation. " "Review and merge are separate workflow roles. A reviewer approval is not merge authorization.") + elif role == "controller": + guidance.append( + "Controller profile: route work via " + "gitea_route_task_session(task_type='process_work_queue') then " + "gitea_allocate_next_work (allocation_mode=cross_role by default). " + "The allocator returns exactly one authoritative selection with " + "required_role / required_profile / selected_action. Do not " + "implement, review, approve, or merge in this session — schedule " + "the matching role namespace instead. Dashboard output is " + "explanatory only and never replaces allocator selection.") elif role == "mixed": guidance.append( "WARNING: this profile allows both authoring and " @@ -15806,7 +15830,8 @@ def gitea_whoami( "environment": profile.get("environment"), "service": profile.get("service"), "identity": profile.get("identity"), - "role": profile.get("role"), + "role": profile.get("role") or _profile_role_kind(profile), + "role_kind": _profile_role_kind(profile), "profile_address": profile.get("profile_path"), "execution_profile": profile.get("execution_profile"), "audit_label": profile.get("audit_label"), @@ -20590,9 +20615,10 @@ def gitea_allocate_next_work( candidates_json: Any = None, exclude_issue_numbers: list[int] | None = None, expected_candidate_set_fingerprint: str | None = None, + allocation_mode: str | None = None, limit: int = 50, ) -> dict: - """Controller-owned next-work allocator using the #613 control-plane DB (#600). + """Controller-owned next-work allocator using the #613 control-plane DB (#600/#840). Workers must not self-select exclusive work under the standard multi-LLM workflow. Call this tool instead. @@ -20602,6 +20628,14 @@ def gitea_allocate_next_work( ``ControlPlaneDB.assign_and_lease`` (never file locks or comment-only leases as the coordination source). + *allocation_mode* (#840): when the active role is controller (or mode is + ``cross_role``), inspect the complete queue and return exactly one + authoritative selection with selected item, action, required_role, + required_profile/namespace, pins, and allocation/lease evidence. + Role-scoped workers pass ``role=author|reviewer|merger|reconciler`` (or + omit for profile role) for single-role filtering. Controller routes only + and does not perform downstream mutations. + Outcomes include: ``assigned_work``, ``preview``, ``wait``, ``blocked_by_terminal_path``, ``no_safe_work``, ``role_ineligible``, ``blocked_by_excluded_own_lease``, ``candidate_set_drift``. @@ -20736,6 +20770,7 @@ def gitea_allocate_next_work( controller_instance_id=allocator_service.resolve_controller_instance_id(), exclude_issue_numbers=exclude_issue_numbers, expected_candidate_set_fingerprint=expected_candidate_set_fingerprint, + allocation_mode=allocation_mode, ) except ValueError as exc: return { diff --git a/namespace_workspace_binding.py b/namespace_workspace_binding.py index fa006b6..be2ca82 100644 --- a/namespace_workspace_binding.py +++ b/namespace_workspace_binding.py @@ -24,7 +24,12 @@ ROLE_WORKTREE_ENVS: dict[str, str] = { "reconciler": RECONCILER_WORKTREE_ENV, } -NON_AUTHOR_ROLES = frozenset({"reviewer", "merger", "reconciler"}) +# Controller has no task worktree env — it routes only (#840). +KNOWN_ROLE_KINDS = frozenset( + {"author", "reviewer", "merger", "reconciler", "controller"} +) + +NON_AUTHOR_ROLES = frozenset({"reviewer", "merger", "reconciler", "controller"}) def normalize_role_kind( @@ -37,8 +42,12 @@ def normalize_role_kind( profile = (profile_name or "").strip().lower() if role == "reviewer" and "merger" in profile: return "merger" + if "controller" in profile or role == "controller": + return "controller" if role in ROLE_WORKTREE_ENVS: return role + if role in KNOWN_ROLE_KINDS: + return role return "author" @@ -80,7 +89,7 @@ def resolve_namespace_workspace( """ env_map = env if env is not None else os.environ role = normalize_role_kind(role_kind, profile_name=profile_name) - role_env_key = ROLE_WORKTREE_ENVS[role] + role_env_key = ROLE_WORKTREE_ENVS.get(role) # #618: durable author resolution — no silent control/master fallback. if role == "author" and verify_paths: @@ -108,13 +117,17 @@ def resolve_namespace_workspace( ) return workspace, source + role_env_candidate = ( + (_env_value(env_map, role_env_key), f"{role_env_key} environment variable", True) + if role_env_key + else (None, "no role worktree env", True) + ) for candidate, source, env_sourced in ( (worktree_path, "worktree_path argument", False), (worktree, "worktree argument", False), (_env_value(env_map, ACTIVE_WORKTREE_ENV), f"{ACTIVE_WORKTREE_ENV} environment variable", True), - (_env_value(env_map, role_env_key), - f"{role_env_key} environment variable", True), + role_env_candidate, (session_lease_worktree if role in {"reviewer", "merger"} else None, "reviewer PR lease worktree", False), # Author lock derivation is handled by the durable path above when @@ -433,7 +446,8 @@ def assess_namespace_mutation_workspace( reasons.append( f"{role} mutation blocked: workspace is the stable control checkout; " f"create or reconnect to a session-owned worktree under branches/ " - f"or set {ROLE_WORKTREE_ENVS[role]} / {ACTIVE_WORKTREE_ENV}" + f"or set {ROLE_WORKTREE_ENVS.get(role, ACTIVE_WORKTREE_ENV)} / " + f"{ACTIVE_WORKTREE_ENV}" ) elif ( role in {"reviewer", "merger"} diff --git a/role_session_router.py b/role_session_router.py index 4f331b8..f756c71 100644 --- a/role_session_router.py +++ b/role_session_router.py @@ -81,6 +81,12 @@ AUTHOR_TASKS = frozenset({ "reconcile_landed_pr", }) +CONTROLLER_TASKS = frozenset({ + "process_work_queue", + "process-work-queue", + "cross_role_allocate", +}) + RECONCILER_TASKS = frozenset({ "cleanup_merged_pr_branch", # #729: delete_branch is reconciler-owned (gitea.branch.delete is granted @@ -132,6 +138,10 @@ TASK_REQUIRED_ROLE = { "reconcile_close_superseded_pr": "reconciler", "reconcile_close_satisfied_issue": "reconciler", "reconcile_create_followup_issue": "reconciler", + # #840: controller-owned generic queue allocation / routing. + "process_work_queue": "controller", + "process-work-queue": "controller", + "cross_role_allocate": "controller", } WRONG_ROLE_REVIEWER_MSG = ( @@ -147,6 +157,10 @@ WRONG_ROLE_MERGER_MSG = ( "Wrong role/session for merger task. Launch merger MCP namespace." ) +WRONG_ROLE_CONTROLLER_MSG = ( + "Wrong role/session for controller task. Launch controller MCP namespace." +) + _session_last_route: dict | None = None @@ -281,6 +295,27 @@ def route_task_session( _record_route(result) return result + if required_role == "controller": + result = { + "task_type": task_type, + "required_role": required_role, + "active_role": active_role_kind, + "active_profile": active_profile, + "route_result": ROUTE_WRONG_ROLE, + "downstream_allowed": False, + "reasons": [ + WRONG_ROLE_CONTROLLER_MSG, + "Controller tasks (process_work_queue / cross-role allocate) " + "cannot run in author, reviewer, merger, or reconciler " + "worker sessions.", + ], + "message": WRONG_ROLE_CONTROLLER_MSG, + "runtime_switching_supported": runtime_switching_supported, + "profile_switch_blocked": not runtime_switching_supported, + } + _record_route(result) + return result + if required_role == "author": route = ROUTE_TO_AUTHOR message = ( diff --git a/task_capability_map.py b/task_capability_map.py index 0a21063..0b8ac0b 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -309,8 +309,10 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "permission": "gitea.pr.create", "role": "author", }, - # #600: controller-owned allocator — any authenticated profile may call; - # routing enforces role match to selected work. Uses control-plane DB (#613). + # #600: workers and controller may call with gitea.read; role-scoped workers + # pass role=author|reviewer|merger|reconciler. Cross-role routing is the + # controller default (#840). The canonical generic queue *task type* is + # process_work_queue (controller-only below). "allocate_next_work": { "permission": "gitea.read", "role": "author", @@ -319,6 +321,19 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "permission": "gitea.read", "role": "author", }, + # #840: documented generic queue task — controller routes only. + "process_work_queue": { + "permission": "gitea.read", + "role": "controller", + }, + "process-work-queue": { + "permission": "gitea.read", + "role": "controller", + }, + "cross_role_allocate": { + "permission": "gitea.read", + "role": "controller", + }, # #601 first-class lease lifecycle — inspect/list need read; mutations gate on # ownership in the control-plane DB (not a separate Gitea write permission). diff --git a/tests/test_cross_role_queue_allocation.py b/tests/test_cross_role_queue_allocation.py new file mode 100644 index 0000000..4126d30 --- /dev/null +++ b/tests/test_cross_role_queue_allocation.py @@ -0,0 +1,581 @@ +"""Authoritative controller cross-role generic queue allocation (#840).""" + +from __future__ import annotations + +import os +import tempfile +import unittest +from unittest.mock import patch + +from allocator_service import ( + ALLOCATION_MODE_CROSS_ROLE, + ALLOCATION_MODE_ROLE_SCOPED, + OUTCOME_NO_SAFE, + OUTCOME_PREVIEW, + OUTCOME_WAIT, + ROLE_AUTHOR, + ROLE_CONTROLLER, + ROLE_MERGER, + ROLE_RECONCILER, + ROLE_REVIEWER, + WorkCandidate, + allocate_next_work, + build_selection_dict, + classify_skip, + required_namespace_for_role, + required_profile_for_role, + resolve_allocation_mode, + selected_action_for_candidate, +) +from control_plane_db import ControlPlaneDB +import role_session_router +from role_session_router import ( + ROUTE_ALLOWED, + ROUTE_AMBIGUOUS, + ROUTE_WRONG_ROLE, + route_task_session, +) +import namespace_workspace_binding as nwb +import task_capability_map + + +class CrossRoleAllocationModeTest(unittest.TestCase): + def test_controller_defaults_to_cross_role(self) -> None: + self.assertEqual( + resolve_allocation_mode(ROLE_CONTROLLER), + ALLOCATION_MODE_CROSS_ROLE, + ) + + def test_worker_defaults_to_role_scoped(self) -> None: + for role in (ROLE_AUTHOR, ROLE_REVIEWER, ROLE_MERGER, ROLE_RECONCILER): + self.assertEqual( + resolve_allocation_mode(role), + ALLOCATION_MODE_ROLE_SCOPED, + ) + + def test_explicit_modes(self) -> None: + self.assertEqual( + resolve_allocation_mode(ROLE_CONTROLLER, "role_scoped"), + ALLOCATION_MODE_ROLE_SCOPED, + ) + self.assertEqual( + resolve_allocation_mode(ROLE_AUTHOR, "cross_role"), + ALLOCATION_MODE_CROSS_ROLE, + ) + + +class CrossRoleSelectionPayloadTest(unittest.TestCase): + def test_selection_contains_required_fields(self) -> None: + c = WorkCandidate( + kind="issue", + number=840, + labels=("status:ready",), + title="cross-role", + priority=20, + ) + sel = build_selection_dict( + c, + active_role=ROLE_CONTROLLER, + required_role=ROLE_AUTHOR, + profile_name="prgs-controller", + allocation_mode=ALLOCATION_MODE_CROSS_ROLE, + ) + self.assertEqual(sel["number"], 840) + self.assertEqual(sel["kind"], "issue") + self.assertEqual(sel["required_role"], ROLE_AUTHOR) + self.assertEqual(sel["selected_action"], "implement") + self.assertEqual(sel["action"], "implement") + self.assertEqual(sel["required_profile"], "prgs-author") + self.assertEqual(sel["required_namespace"], "gitea-author") + self.assertEqual(sel["pinned"]["number"], 840) + self.assertIsNone(sel["pinned"]["head_sha"]) + + def test_profile_prefix_preserved(self) -> None: + self.assertEqual( + required_profile_for_role(ROLE_REVIEWER, profile_name="dadeschools-controller"), + "dadeschools-reviewer", + ) + self.assertEqual( + required_namespace_for_role(ROLE_MERGER), + "gitea-merger", + ) + + def test_selected_actions_per_role(self) -> None: + issue = WorkCandidate(kind="issue", number=1, labels=("status:ready",)) + pr_review = WorkCandidate(kind="pr", number=2, head_sha="a" * 40) + pr_rc = WorkCandidate( + kind="pr", + number=3, + head_sha="b" * 40, + request_changes_current_head=True, + ) + pr_merge = WorkCandidate( + kind="pr", + number=4, + head_sha="c" * 40, + approval_on_current_head=True, + mergeable=True, + ) + pr_recon = WorkCandidate( + kind="pr", + number=5, + head_sha="d" * 40, + approval_contaminated=True, + ) + self.assertEqual(selected_action_for_candidate(issue, ROLE_AUTHOR), "implement") + self.assertEqual( + selected_action_for_candidate(pr_rc, ROLE_AUTHOR), + "address_pr_change_requests", + ) + self.assertEqual( + selected_action_for_candidate(pr_review, ROLE_REVIEWER), "review" + ) + self.assertEqual(selected_action_for_candidate(pr_merge, ROLE_MERGER), "merge") + self.assertEqual( + selected_action_for_candidate(pr_recon, ROLE_RECONCILER), + "reconcile_contaminated_approval", + ) + + +class CrossRoleAllocateServiceTest(unittest.TestCase): + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory() + self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3")) + + def tearDown(self) -> None: + self._tmp.cleanup() + + def _alloc(self, **kwargs): + defaults = dict( + db=self.db, + session_id="ctrl-session", + role=ROLE_CONTROLLER, + remote="prgs", + org="org", + repo="repo", + candidates=[], + apply=False, + profile_name="prgs-controller", + username="controller-bot", + controller_instance_id="ctrl-1", + ) + defaults.update(kwargs) + return allocate_next_work(**defaults) + + def test_eligible_author_work(self) -> None: + cands = [ + WorkCandidate( + kind="issue", + number=100, + labels=("status:ready",), + title="author work", + priority=20, + ), + ] + res = self._alloc(candidates=cands) + self.assertTrue(res["success"]) + self.assertEqual(res["outcome"], OUTCOME_PREVIEW) + self.assertEqual(res["allocation_mode"], ALLOCATION_MODE_CROSS_ROLE) + self.assertIsNotNone(res["selected"]) + self.assertEqual(res["selected"]["number"], 100) + self.assertEqual(res["required_role"], ROLE_AUTHOR) + self.assertEqual(res["selected_action"], "implement") + self.assertEqual(res["required_profile"], "prgs-author") + self.assertEqual(res["required_namespace"], "gitea-author") + self.assertIn("allocate", res["controller_allowed_actions"]) + self.assertIn("merge", res["controller_forbidden_actions"]) + self.assertFalse(res["allocation_evidence"]["lease_created"]) + + def test_eligible_reviewer_work(self) -> None: + cands = [ + WorkCandidate( + kind="pr", + number=200, + head_sha="e" * 40, + title="needs review", + priority=30, + ), + ] + res = self._alloc(candidates=cands) + self.assertEqual(res["selected"]["number"], 200) + self.assertEqual(res["required_role"], ROLE_REVIEWER) + self.assertEqual(res["selected_action"], "review") + self.assertEqual(res["required_profile"], "prgs-reviewer") + self.assertEqual(res["selected"]["pinned"]["head_sha"], "e" * 40) + + def test_eligible_merger_work(self) -> None: + cands = [ + WorkCandidate( + kind="pr", + number=300, + head_sha="f" * 40, + approval_on_current_head=True, + mergeable=True, + priority=40, + ), + ] + res = self._alloc(candidates=cands) + self.assertEqual(res["selected"]["number"], 300) + self.assertEqual(res["required_role"], ROLE_MERGER) + self.assertEqual(res["selected_action"], "merge") + + def test_eligible_reconciler_work(self) -> None: + cands = [ + WorkCandidate( + kind="pr", + number=400, + head_sha="1" * 40, + approval_contaminated=True, + priority=50, + ), + ] + res = self._alloc(candidates=cands) + self.assertEqual(res["selected"]["number"], 400) + self.assertEqual(res["required_role"], ROLE_RECONCILER) + self.assertIn("reconcile", res["selected_action"]) + + def test_no_eligible_work(self) -> None: + cands = [ + WorkCandidate( + kind="issue", + number=10, + labels=("status:blocked",), + blocked=True, + priority=99, + ), + WorkCandidate( + kind="issue", + number=11, + labels=("status:ready",), + dependency_unmet=True, + dependency_reason="blocked by #10", + priority=98, + ), + ] + res = self._alloc(candidates=cands) + self.assertTrue(res["success"]) + self.assertEqual(res["outcome"], OUTCOME_NO_SAFE) + self.assertIsNone(res["selected"]) + self.assertEqual(res["allocation_mode"], ALLOCATION_MODE_CROSS_ROLE) + + def test_leased_work_skipped(self) -> None: + cands = [ + WorkCandidate( + kind="issue", + number=50, + labels=("status:ready",), + priority=20, + ), + WorkCandidate( + kind="issue", + number=51, + labels=("status:ready",), + priority=10, + ), + ] + # Seed a foreign lease on issue 50 via assign_and_lease under another session. + other = allocate_next_work( + self.db, + session_id="other-worker", + role=ROLE_AUTHOR, + remote="prgs", + org="org", + repo="repo", + candidates=cands[:1], + apply=True, + profile_name="prgs-author", + controller_instance_id="other-ctrl", + ) + self.assertEqual(other["outcome"], "assigned_work") + res = self._alloc(candidates=cands) + self.assertIsNotNone(res["selected"]) + self.assertEqual(res["selected"]["number"], 51) + self.assertTrue(any(s["number"] == 50 for s in res["skipped"])) + self.assertTrue(res["claims_excluded"]) + + def test_dependencies_skipped(self) -> None: + cands = [ + WorkCandidate( + kind="issue", + number=1, + labels=("status:ready",), + priority=99, + dependency_unmet=True, + dependency_reason="needs #2", + ), + WorkCandidate( + kind="issue", + number=2, + labels=("status:ready",), + priority=1, + ), + ] + res = self._alloc(candidates=cands) + self.assertEqual(res["selected"]["number"], 2) + skipped = {s["number"]: s["reason"] for s in res["skipped"]} + self.assertIn(1, skipped) + self.assertIn("needs #2", skipped[1]) + + def test_pagination_limit_only_truncates_skip_report(self) -> None: + """Ranking uses full inventory; reporting limit is MCP-layer only. + + Service ranks all candidates; prove higher-priority eligible item + wins even when many skipped precede it. + """ + cands = [] + for n in range(1, 30): + cands.append( + WorkCandidate( + kind="issue", + number=n, + labels=("status:ready",), + priority=100 - n, + dependency_unmet=True, + dependency_reason=f"dep {n}", + ) + ) + cands.append( + WorkCandidate( + kind="issue", + number=999, + labels=("status:ready",), + priority=1, + ) + ) + res = self._alloc(candidates=cands) + self.assertEqual(res["selected"]["number"], 999) + self.assertGreaterEqual(len(res["skipped"]), 29) + + def test_role_scoped_controller_legacy_still_restricts(self) -> None: + """role_scoped controller only takes reconciler-needed items.""" + cands = [ + WorkCandidate( + kind="issue", + number=1, + labels=("status:ready",), + priority=50, + ), + WorkCandidate( + kind="pr", + number=2, + head_sha="a" * 40, + approval_contaminated=True, + priority=1, + ), + ] + res = self._alloc( + candidates=cands, + allocation_mode=ALLOCATION_MODE_ROLE_SCOPED, + ) + self.assertEqual(res["allocation_mode"], ALLOCATION_MODE_ROLE_SCOPED) + self.assertEqual(res["selected"]["number"], 2) + self.assertEqual(res["required_role"], ROLE_RECONCILER) + + def test_cross_role_prefers_highest_priority_across_roles(self) -> None: + cands = [ + WorkCandidate( + kind="issue", + number=10, + labels=("status:ready",), + priority=10, + ), + WorkCandidate( + kind="pr", + number=20, + head_sha="b" * 40, + priority=50, + ), + WorkCandidate( + kind="pr", + number=30, + head_sha="c" * 40, + approval_on_current_head=True, + mergeable=True, + priority=20, + ), + ] + res = self._alloc(candidates=cands) + # PR #20 highest priority → reviewer + self.assertEqual(res["selected"]["number"], 20) + self.assertEqual(res["required_role"], ROLE_REVIEWER) + + def test_apply_creates_lease_evidence_for_required_role(self) -> None: + cands = [ + WorkCandidate( + kind="issue", + number=777, + labels=("status:ready",), + priority=20, + ), + ] + res = self._alloc(candidates=cands, apply=True) + self.assertEqual(res["outcome"], "assigned_work") + self.assertTrue(res["allocation_evidence"]["lease_created"]) + self.assertEqual(res["allocation_evidence"]["lease_role"], ROLE_AUTHOR) + proof = res["lease_proof"] + self.assertIsNotNone(proof["lease_id"]) + self.assertEqual(proof["lease_role"], ROLE_AUTHOR) + self.assertIn("implement", proof["allowed_actions"]) + # Controller isolation: controller still forbids merge/push/create_pr + self.assertIn("merge", res["controller_forbidden_actions"]) + self.assertIn("push", res["controller_forbidden_actions"]) + + def test_metadata_consistency_role_is_controller(self) -> None: + cands = [ + WorkCandidate( + kind="issue", + number=1, + labels=("status:ready",), + ), + ] + res = self._alloc(candidates=cands) + self.assertEqual(res["role"], ROLE_CONTROLLER) + self.assertEqual(res["routing_role"], ROLE_CONTROLLER) + self.assertEqual(res["required_role"], ROLE_AUTHOR) + + +class ProcessWorkQueueRouterTest(unittest.TestCase): + def tearDown(self) -> None: + role_session_router.clear_route_state() + + def test_process_work_queue_allowed_for_controller(self) -> None: + res = route_task_session( + "process_work_queue", + active_profile="prgs-controller", + active_role_kind="controller", + allowed_in_current_session=True, + ) + self.assertEqual(res["route_result"], ROUTE_ALLOWED) + self.assertEqual(res["required_role"], "controller") + self.assertTrue(res["downstream_allowed"]) + + def test_process_work_queue_hyphen_alias(self) -> None: + res = route_task_session( + "process-work-queue", + active_profile="prgs-controller", + active_role_kind="controller", + allowed_in_current_session=True, + ) + self.assertEqual(res["route_result"], ROUTE_ALLOWED) + + def test_process_work_queue_wrong_role_for_author(self) -> None: + res = route_task_session( + "process_work_queue", + active_profile="prgs-author", + active_role_kind="author", + allowed_in_current_session=False, + ) + self.assertEqual(res["route_result"], ROUTE_WRONG_ROLE) + self.assertEqual(res["required_role"], "controller") + self.assertFalse(res["downstream_allowed"]) + + def test_unknown_still_ambiguous(self) -> None: + res = route_task_session( + "not_a_real_task", + active_profile="prgs-controller", + active_role_kind="controller", + allowed_in_current_session=False, + ) + self.assertEqual(res["route_result"], ROUTE_AMBIGUOUS) + + def test_capability_map_process_work_queue_is_controller(self) -> None: + self.assertEqual( + task_capability_map.required_role("process_work_queue"), + "controller", + ) + self.assertEqual( + task_capability_map.required_permission("process_work_queue"), + "gitea.read", + ) + + +class ControllerRoleMetadataTest(unittest.TestCase): + def test_normalize_role_kind_controller(self) -> None: + self.assertEqual( + nwb.normalize_role_kind("controller"), + "controller", + ) + self.assertEqual( + nwb.normalize_role_kind("author", profile_name="prgs-controller"), + "controller", + ) + self.assertEqual( + nwb.normalize_role_kind("reconciler", profile_name="prgs-controller"), + "controller", + ) + + def test_profile_role_kind_prefers_declared_controller(self) -> None: + # Import from worktree package path via sys.path already set by pytest. + import gitea_mcp_server as mcp + + profile = { + "profile_name": "prgs-controller", + "role": "controller", + "allowed_operations": [ + "gitea.read", + "gitea.issue.comment", + "gitea.pr.close", + ], + "forbidden_operations": [ + "gitea.pr.approve", + "gitea.pr.merge", + "gitea.pr.create", + "gitea.branch.push", + ], + } + # Declared role wins even if permissions look reconciler-like. + self.assertEqual(mcp._profile_role_kind(profile), "controller") + # Name-based fallback. + profile_no_role = dict(profile) + profile_no_role["role"] = None + profile_no_role["role_kind"] = None + self.assertEqual(mcp._profile_role_kind(profile_no_role), "controller") + + def test_permission_inference_without_controller_name_stays_reconciler(self) -> None: + import gitea_mcp_server as mcp + + # Pure permission inference still may return reconciler when no controller + # declaration exists — that is intentional for reconciler profiles. + role = mcp._role_kind( + ["gitea.read", "gitea.pr.close", "gitea.issue.comment"], + ["gitea.pr.approve", "gitea.pr.merge", "gitea.pr.create", "gitea.branch.push"], + ) + self.assertEqual(role, "reconciler") + + +class DashboardRemainsExplanatoryTest(unittest.TestCase): + def test_dashboard_prompt_points_at_allocator_not_self_select(self) -> None: + import workflow_dashboard as wd + + self.assertIn("gitea_allocate_next_work", wd.PROMPT_CONTROLLER) + self.assertIn("process_work_queue", wd.PROMPT_CONTROLLER) + self.assertIn("never replaces allocator", wd.PROMPT_CONTROLLER.lower()) + self.assertNotIn("self-select", wd.PROMPT_CONTROLLER.lower()) + + +class ClassifySkipCrossRoleTest(unittest.TestCase): + def test_controller_cross_role_accepts_author_issue(self) -> None: + c = WorkCandidate(kind="issue", number=1, labels=("status:ready",)) + self.assertIsNone( + classify_skip( + c, + role=ROLE_CONTROLLER, + terminal_pr=None, + allocation_mode=ALLOCATION_MODE_CROSS_ROLE, + ) + ) + + def test_legacy_controller_skips_author_issue(self) -> None: + c = WorkCandidate(kind="issue", number=1, labels=("status:ready",)) + reason = classify_skip( + c, + role=ROLE_CONTROLLER, + terminal_pr=None, + allocation_mode=ALLOCATION_MODE_ROLE_SCOPED, + ) + self.assertIsNotNone(reason) + self.assertIn("does not require controller", reason or "") + + +if __name__ == "__main__": + unittest.main() diff --git a/workflow_dashboard.py b/workflow_dashboard.py index d2b10ee..e47b2dc 100644 --- a/workflow_dashboard.py +++ b/workflow_dashboard.py @@ -65,9 +65,11 @@ PROMPT_RECONCILER = ( "reconciliation (already-landed / post-merge cleanup). Do not approve or merge." ) PROMPT_CONTROLLER = ( - "CONTROLLER session: inspect gitea_workflow_dashboard + control-plane leases, " - "diagnose blocked/terminal-locked items for {remote}/{org}/{repo}, and schedule " - "exactly one fresh role-scoped cycle. Do not implement, review, or merge in-band." + "CONTROLLER session: call gitea_route_task_session(task_type='process_work_queue') " + "then gitea_allocate_next_work (cross_role default) for {remote}/{org}/{repo}; " + "use the returned required_role/profile/action to schedule exactly one downstream " + "role cycle. Dashboard is explanatory only and never replaces allocator selection. " + "Do not implement, review, approve, or merge in-band." ) PROMPT_IDLE = ( "IDLE: no safe assignable work for role '{role}' on {remote}/{org}/{repo}. "