Compare commits

..
Author SHA1 Message Date
sysadmin a20975688d Merge branch 'master' into feat/issue-637-timeline-model
# Conflicts:
#	webui/app.py
2026-07-23 14:53:12 -04:00
sysadminandClaude Opus 4.8 25bc2a3291 feat(webui): workflow-event and conversation timeline model (Closes #637)
Phase 1 child of the Web Console epic #631. Adds a durable, versioned
WorkflowEvent schema with per-source adapters and a read-only query API so
operators can browse a unified timeline of workflow events, decisions, tool
calls, and handoffs instead of scattered evidence.

- webui/timeline.py (new): versioned WorkflowEvent schema; control-plane
  event adapter and Gitea CTH handoff-comment adapter; read-only mode=ro
  control-plane reader; conjunctive filter by issue/PR/session; stable
  (timestamp, source_rank, event_key) ordering; bounded pagination;
  fail-soft per-source status; redaction at the boundary, fail closed.
- webui/app.py: GET /api/v1/timeline read-only route with thread-scoped,
  fail-soft handoff comment source.
- tests/test_webui_timeline.py (new): schema, adapters, redaction of
  secret-like payloads, filter/sort/pagination, scoped CP reader,
  fail-soft composition, and API integration.
- docs/webui-local-dev.md: timeline route and field-authority notes.

Read-only Phase 1; no mutation of historical events; no full chat replay;
no unredacted tool-argument storage.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-23 14:48:46 -04:00
sysadmin 1c455b6ec0 Merge pull request 'feat(webui): read-only system-health API (Closes #634)' (#813) from feat/issue-634-readonly-system-health-api into master 2026-07-23 04:14:33 -05:00
jcwalker3 6868b345ee Merge branch 'master' into feat/issue-634-readonly-system-health-api 2026-07-23 01:12:52 -05:00
jcwalker3 da6a864463 Merge branch 'master' into feat/issue-634-readonly-system-health-api 2026-07-23 00:06:06 -05:00
jcwalker3andClaude Opus 4.8 5494696227 feat(webui): read-only system-health API (Closes #634)
Adds `GET /api/v1/system/health`, a structured read-only health surface for
automated readiness checks, and keeps `/health` as the cheap liveness probe.

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

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

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

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

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

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-22 15:59:30 -05:00
10 changed files with 2366 additions and 746 deletions
+151 -16
View File
@@ -52,7 +52,8 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| 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 with status and onboarding progress (#427, #635) |
@@ -64,8 +65,6 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| `/api/prompts` | JSON prompt export with workflow hashes |
| `/runtime` | MCP runtime health and stale detection (#430) |
| `/api/runtime` | JSON runtime health export |
| `/policy` | Workflow policy and guardrail configuration visibility (#646) |
| `/api/v1/policy` | Versioned JSON guardrail inventory (redacted, read-only) |
| `/audit` | Report audit paste + validator preview (#431) |
| `/api/audit` | JSON validator preview (POST `report_text`, optional `task_kind`) |
| `/worktrees` | Worktree hygiene dashboard (#432) |
@@ -80,6 +79,85 @@ Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
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
@@ -155,19 +233,6 @@ health, workflow/schema SHA-256 hashes, and stale-runtime warnings when the
checkout is behind merged safety-gate changes. Restart guidance links to #420;
no tokens or MCP restart actions are exposed.
## Policy & guardrail visibility (#646)
`/policy` (HTML) and `/api/v1/policy` (JSON) surface a **read-only** projection
of the major workflow guardrails — role separation/RBAC, lease lifecycle,
author worktree binding, merge confirmation, secret redaction, contamination
containment, allocator policy, audit logging, and mutation gating. Each entry
carries source pointers to the file/module/doc that owns it, a compact active
value derived from the existing safe policy accessors, and — where a documented
default is declared — a diff of active vs documented. The whole payload is run
through the console redaction pass before it is emitted, so a planted or
accidental secret degrades to the placeholder rather than reaching a client.
The view never edits policy and exposes no gate-weakening toggle.
## Deployment boundary (#435)
MVP serves on loopback by default. Binding `0.0.0.0` or `::` is **refused**
@@ -227,6 +292,76 @@ health, workflow/schema SHA-256 hashes, and stale-runtime warnings when the
checkout is behind merged safety-gate changes. Restart guidance links to #420;
no tokens or MCP restart actions are exposed.
## Workflow-event timeline (#637)
`GET /api/v1/timeline` is a read-only, versioned aggregation of workflow
events from every available source into one normalised, filterable stream. It
is the model layer for the Phase 1 timeline console view (a later child issue
of #631); this issue ships the schema, adapters, and read API only.
### Schema (versioned)
`webui/timeline.py` declares `TIMELINE_SCHEMA_VERSION` (currently `1`) and the
frozen `WorkflowEvent` record. Every response carries `schema_version` so a
consumer can branch on shape. One event:
```json
{
"source": "control_plane",
"event_type": "lease.renew",
"event_key": "cp:1421",
"timestamp": "2026-07-23T02:00:00Z",
"actor": null,
"role": null,
"issue_number": 637,
"pr_number": null,
"session_id": null,
"tool_name": null,
"decision": null,
"message": "lease renewed",
"correlation_id": "issue#637",
"evidence_refs": [],
"sensitive": true
}
```
`event_key` is stable and unique per source (`cp:<event_id>`,
`cth:<kind>:<number>:<comment_id>`), so pagination and dedup are deterministic.
### Sources and field authority
| Source | Adapter | Authority |
|---|---|---|
| Control-plane `events``work_items` | `adapt_cp_events` | `event_type`, `message`, `timestamp`, issue/PR scope come from the CP database, read through a `mode=ro` URI (never creates the DB or runs migrations) |
| Gitea Canonical Thread Handoff comments | `adapt_cth_comments` | `actor`, `role` (next owner), `decision`, `evidence_refs`, `timestamp` come from the parsed CTH comment body (`canonical_thread_handoff`) |
Handoff comments are thread-scoped: they are only read when the request filters
by a single `issue` or `pr`. Otherwise the handoff source reports `not run`
with a reason — it is never rendered as empty-and-healthy. Each source degrades
independently: an unavailable control-plane DB or a failed comment fetch is a
`sources[]` entry with `ok:false` and a `reason`, never a dropped timeline.
### Query parameters
`issue`, `pr`, `session` (conjunctive filters); `limit` (default 50, max 500)
and `offset` for pagination; `remote`, `org`, `repo` to override the default
registry-project scope. Events sort ascending by
`(timestamp, source_rank, event_key)`; missing timestamps sort last.
### Redaction
Every free-text field (event messages, decision/proof text, roles) is passed
through the console redaction policy (`webui.console_redaction`, backed by
`gitea_audit.redact`) before it leaves the module, failing closed to the
placeholder. No unredacted tool arguments or secrets are ever emitted, and a
generation error never drops raw data to a caller or a log.
### Tests
```bash
pytest tests/test_webui_timeline.py -q
```
## Tests
```bash
-222
View File
@@ -1,222 +0,0 @@
"""Tests for the read-only workflow policy/guardrail visibility view (#646).
Covers issue #646 acceptance criteria:
1. Console lists major guardrails with source pointers.
2. Secrets redacted.
3. Tests ensure sample secrets never appear.
4. Docs explain read-only nature (asserted here for the page copy; the doc
itself is covered by inspection).
"""
import json
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient
from webui import console_redaction
from webui import policy_inventory
from webui.app import create_app
from webui.policy_inventory import (
PolicyEntry,
PolicyInventorySnapshot,
SourcePointer,
load_policy_inventory,
snapshot_to_dict,
)
from webui.policy_views import render_policy_page
def _entry(key, category, *, active=None, error=None):
return PolicyEntry(
key=key,
title=key.replace("_", " ").title(),
category=category,
summary=f"summary for {key}",
sources=(SourcePointer("src", f"{key}.py", "module"),),
active=active,
documented_default=None,
diff=None,
error=error,
)
def _snapshot(entries):
return PolicyInventorySnapshot(
schema_version=1,
read_only=True,
note="read-only projection",
entries=tuple(entries),
categories=tuple(dict.fromkeys(e.category for e in entries)),
build_errors=(),
)
# The guardrail categories issue #646 names as in-scope.
_EXPECTED_CATEGORIES = {
"role_separation",
"lease_rules",
"worktree_rules",
"merge_confirmation",
"redaction",
"contamination",
"allocator_policy",
"audit_logging",
"mutation_gating",
}
class TestPolicyInventoryModel(unittest.TestCase):
def test_major_guardrails_present(self):
snapshot = load_policy_inventory()
categories = {e.category for e in snapshot.entries}
self.assertEqual(_EXPECTED_CATEGORIES, categories)
self.assertGreaterEqual(len(snapshot.entries), len(_EXPECTED_CATEGORIES))
def test_every_guardrail_has_source_pointers(self):
# AC1: source attribution (file/module/doc) for every guardrail.
snapshot = load_policy_inventory()
for entry in snapshot.entries:
with self.subTest(entry=entry.key):
self.assertTrue(entry.sources, "guardrail must carry source pointers")
for source in entry.sources:
self.assertTrue(source.path)
self.assertIn(source.kind, {"module", "doc", "script", "config"})
def test_diff_reported_where_documented_default_declared(self):
snapshot = load_policy_inventory()
checked_any = False
for entry in snapshot.entries:
if entry.documented_default is None:
self.assertIsNone(entry.diff)
continue
checked_any = True
self.assertIsNotNone(entry.diff)
self.assertEqual(
entry.diff["status"],
"matches_documented_default",
f"{entry.key} drifted from its documented default: {entry.diff}",
)
self.assertTrue(checked_any, "at least one guardrail should declare a default")
def test_live_projections_populate_active(self):
snapshot = load_policy_inventory()
by_key = {e.key: e for e in snapshot.entries}
for key in ("role_separation", "redaction", "audit_logging"):
self.assertIsNone(by_key[key].error, f"{key} projection failed")
self.assertIsInstance(by_key[key].active, dict)
def test_build_entry_is_fail_soft_on_projection_error(self):
def _boom():
raise RuntimeError("projection exploded")
row = (
"redaction",
"Secret redaction",
"redaction",
"summary",
(SourcePointer("x", "webui/console_redaction.py", "module"),),
_boom,
{"redact_before_persist": True},
)
entry = policy_inventory._build_entry(row)
self.assertIsNone(entry.active)
self.assertIsNotNone(entry.error)
self.assertEqual(entry.diff["status"], "active_unavailable")
class TestPolicyRedaction(unittest.TestCase):
def test_real_snapshot_has_no_secret_shapes(self):
# AC3: the real emitted payload never carries a known secret shape.
payload = snapshot_to_dict(load_policy_inventory())
self.assertEqual(console_redaction.scan_for_secrets(payload), [])
def test_planted_keychain_secret_is_redacted(self):
# AC2/AC3: a secret planted in an active projection is masked before emit.
snapshot = _snapshot([
_entry(
"redaction",
"redaction",
active={"leaked": "keychain:prgs-author-super-secret", "roles": ["author"]},
)
])
payload = snapshot_to_dict(snapshot)
blob = json.dumps(payload)
self.assertNotIn("keychain:prgs-author-super-secret", blob)
self.assertEqual(console_redaction.scan_for_secrets(payload), [])
def test_planted_credential_assignment_is_redacted(self):
snapshot = _snapshot([
_entry(
"audit_logging",
"audit_logging",
active={"leaked": "token=abcd1234efgh5678", "append_only": True},
)
])
payload = snapshot_to_dict(snapshot)
blob = json.dumps(payload)
self.assertNotIn("abcd1234efgh5678", blob)
self.assertEqual(console_redaction.scan_for_secrets(payload), [])
class TestPolicyRoutes(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_policy_html_lists_guardrails_with_sources(self):
response = self.client.get("/policy")
self.assertEqual(response.status_code, 200)
text = response.text
self.assertIn("Workflow policy", text)
self.assertIn("Role separation and RBAC", text)
self.assertIn("Source pointers", text)
self.assertIn("task_capability_map.py", text)
self.assertIn("docs/safety-model.md", text)
def test_policy_html_states_read_only(self):
# AC4: the page explains its read-only nature.
text = self.client.get("/policy").text
self.assertIn("read-only", text.lower())
self.assertNotIn("<form", text.lower())
def test_policy_html_has_no_secret_shapes(self):
text = self.client.get("/policy").text
self.assertEqual(console_redaction.scan_for_secrets(text), [])
def test_api_v1_policy_returns_inventory(self):
response = self.client.get("/api/v1/policy")
self.assertEqual(response.status_code, 200)
data = response.json()
self.assertEqual(data["schema_version"], policy_inventory.SCHEMA_VERSION)
self.assertTrue(data["read_only"])
self.assertEqual(data["entry_count"], len(data["entries"]))
self.assertEqual(set(data["categories"]), _EXPECTED_CATEGORIES)
def test_policy_is_read_only_no_post(self):
# AC4 / non-goal: no mutation endpoint.
response = self.client.post("/policy")
self.assertIn(response.status_code, (404, 405))
def test_nav_links_policy(self):
text = self.client.get("/").text
self.assertIn('href="/policy"', text)
class TestPolicyViewFailSoft(unittest.TestCase):
def test_page_renders_when_a_projection_errors(self):
snapshot = _snapshot([
_entry("role_separation", "role_separation", error="active projection unavailable: boom"),
_entry("redaction", "redaction", active={"redact_before_persist": True}),
])
page = render_policy_page(snapshot)
# The errored guardrail surfaces its error; other guardrails still render.
self.assertIn("Active value unavailable", page)
self.assertIn("Redaction", page)
self.assertIn("Workflow policy", page)
if __name__ == "__main__":
unittest.main()
+499
View File
@@ -0,0 +1,499 @@
"""Tests for the read-only system-health API (#634).
Covers the acceptance criteria directly: a structured payload with readiness
and a dependency list (AC1), version and uptime when knowable (AC2), stale
runtime reported without a false mutation-safe claim (AC3), and the healthy /
degraded-dependency / redaction cases (AC4).
"""
import json
import os
import sqlite3
import sys
import tempfile
import unittest
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient
import control_plane_db
from webui.app import create_app
from webui.deployment_boundary import scan_text_for_client_secrets
from webui.system_health import (
API_PATH,
STATUS_DEGRADED,
STATUS_DOWN,
STATUS_OK,
STATUS_SKIPPED,
DependencyProbe,
StaleRuntime,
assess_stale_runtime,
clear_probe_cache,
load_system_health,
namespace_summaries,
probe_control_plane_db,
probe_gitea,
process_uptime,
redact,
redact_url,
snapshot_to_dict,
)
def _probe(name, status, *, required=True, detail="detail", kind="test"):
return DependencyProbe(
name=name,
kind=kind,
status=status,
detail=detail,
required=required,
latency_ms=1.5,
metadata={},
)
_ALL_HEALTHY = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_OK, required=False, kind="http"),
)
_CLEAN_PARITY = StaleRuntime(
daemon_head="abc123",
checkout_head="abc123",
remote_head="abc123",
stale=False,
determinable=True,
mutation_safe=True,
reasons=(),
)
class CleanParityMixin:
"""Pin parity for tests about aggregation rather than staleness.
Without this the assertions depend on the real checkout: a worktree whose
branch is ahead of its upstream is genuinely stale, which would degrade the
overall status and make these cases fail for an unrelated reason.
"""
def setUp(self):
super().setUp()
patcher = mock.patch(
"webui.system_health.assess_stale_runtime",
return_value=_CLEAN_PARITY,
)
patcher.start()
self.addCleanup(patcher.stop)
class TestDependencyAggregation(CleanParityMixin, unittest.TestCase):
"""AC1 — readiness and dependency list derived from probe results."""
def test_all_healthy_is_ok_and_ready(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc123")
self.assertEqual(snapshot.status, STATUS_OK)
self.assertTrue(snapshot.ready)
self.assertTrue(snapshot.readiness_complete)
self.assertEqual(snapshot.readiness_reasons, ())
self.assertEqual(len(snapshot.dependencies), 3)
def test_required_dependency_down_blocks_readiness(self):
probes = (
_probe("control_plane_db", STATUS_DOWN, detail="file missing", kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_OK, required=False, kind="http"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertEqual(snapshot.status, STATUS_DOWN)
self.assertFalse(snapshot.ready)
self.assertTrue(
any("control_plane_db" in reason for reason in snapshot.readiness_reasons)
)
def test_optional_dependency_down_degrades_but_stays_ready(self):
"""A failing optional probe must not claim the process itself is unready."""
probes = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_DOWN, required=False, detail="timeout", kind="http"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertEqual(snapshot.status, STATUS_DEGRADED)
self.assertTrue(snapshot.ready)
self.assertTrue(any("gitea" in reason for reason in snapshot.readiness_reasons))
def test_unrun_required_probe_leaves_readiness_incomplete(self):
"""Not probed is not the same as passing."""
probes = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_SKIPPED, detail="offline", kind="git"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertFalse(snapshot.ready)
self.assertFalse(snapshot.readiness_complete)
self.assertEqual(snapshot.status, STATUS_DEGRADED)
def test_skipped_optional_probe_does_not_block_readiness(self):
probes = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_SKIPPED, required=False, kind="http"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertTrue(snapshot.ready)
self.assertTrue(snapshot.readiness_complete)
class TestVersionAndUptime(CleanParityMixin, unittest.TestCase):
"""AC2 — version and uptime present when knowable."""
def test_uptime_and_start_time_present(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc123")
self.assertGreaterEqual(snapshot.uptime_seconds, 0.0)
self.assertIn("T", snapshot.started_at)
def test_process_uptime_helper_matches_shape(self):
started_at, uptime = process_uptime()
self.assertIn("T", started_at)
self.assertGreaterEqual(uptime, 0.0)
def test_version_reports_python_and_schema_version(self):
probes = (
DependencyProbe(
name="control_plane_db",
kind="sqlite",
status=STATUS_OK,
detail="ok",
required=True,
latency_ms=1.0,
metadata={"schema_version": control_plane_db.SCHEMA_VERSION},
),
_probe("repository", STATUS_OK, kind="git"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertEqual(
snapshot.version.control_plane_schema_version,
control_plane_db.SCHEMA_VERSION,
)
self.assertTrue(snapshot.version.python_version)
def test_version_known_flag_false_when_sha_unavailable(self):
with mock.patch("webui.system_health._git", return_value=None):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc")
self.assertIsNone(snapshot.version.git_sha)
self.assertFalse(snapshot.version.known)
class TestStaleRuntime(unittest.TestCase):
"""AC3 — stale runtime reflected without a false mutation-safe claim."""
def test_matching_commits_are_mutation_safe(self):
assessment = assess_stale_runtime(
Path("/tmp"),
daemon_head="aaa",
git_reader=lambda *args: "aaa",
)
self.assertFalse(assessment.stale)
self.assertTrue(assessment.determinable)
self.assertTrue(assessment.mutation_safe)
def test_diverged_commits_are_stale_and_not_mutation_safe(self):
reads = {"HEAD": "aaa", "@{upstream}": "bbb"}
assessment = assess_stale_runtime(
Path("/tmp"),
daemon_head="aaa",
git_reader=lambda *args: reads.get(args[-1]),
)
self.assertTrue(assessment.stale)
self.assertFalse(assessment.mutation_safe)
self.assertTrue(assessment.reasons)
def test_unknown_remote_is_not_mutation_safe(self):
"""Indeterminate must never read as safe."""
reads = {"HEAD": "aaa", "@{upstream}": None}
assessment = assess_stale_runtime(
Path("/tmp"),
daemon_head="aaa",
git_reader=lambda *args: reads.get(args[-1]),
)
self.assertFalse(assessment.determinable)
self.assertFalse(assessment.mutation_safe)
self.assertFalse(assessment.stale)
self.assertTrue(
any("indeterminate" in reason for reason in assessment.reasons)
)
def test_unobservable_daemon_head_is_disclosed(self):
assessment = assess_stale_runtime(
Path("/tmp"),
git_reader=lambda *args: "aaa",
)
self.assertTrue(
any("not observable" in reason for reason in assessment.reasons)
)
def test_stale_runtime_degrades_overall_status(self):
reads = {"HEAD": "aaa", "@{upstream}": "bbb"}
# Pinned rather than inherited: this path uses the default git reader,
# so the assertion must hold whether or not the suite runs offline.
with mock.patch.dict(os.environ, {"WEBUI_TEST_OFFLINE": ""}), mock.patch(
"webui.system_health._git",
side_effect=lambda repo, *args: reads.get(args[-1]),
):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="aaa")
self.assertTrue(snapshot.stale_runtime.stale)
self.assertFalse(snapshot.stale_runtime.mutation_safe)
self.assertEqual(snapshot.status, STATUS_DEGRADED)
class TestControlPlaneDbProbe(unittest.TestCase):
"""The required local dependency, probed read-only."""
def setUp(self):
self.tmp = tempfile.TemporaryDirectory()
self.addCleanup(self.tmp.cleanup)
self.db_path = str(Path(self.tmp.name) / "control-plane.db")
def _build_db(self, schema_version):
conn = sqlite3.connect(self.db_path)
conn.execute("CREATE TABLE schema_meta (key TEXT PRIMARY KEY, value TEXT)")
conn.execute("CREATE TABLE leases (lease_id TEXT PRIMARY KEY, status TEXT)")
conn.execute(
"INSERT INTO schema_meta(key, value) VALUES ('schema_version', ?)",
(str(schema_version),),
)
conn.execute("INSERT INTO leases(lease_id, status) VALUES ('l1', 'active')")
conn.commit()
conn.close()
def test_missing_database_is_down(self):
probe = probe_control_plane_db(str(Path(self.tmp.name) / "absent.db"))
self.assertEqual(probe.status, STATUS_DOWN)
self.assertTrue(probe.required)
self.assertIsNotNone(probe.latency_ms)
def test_matching_schema_is_ok(self):
self._build_db(control_plane_db.SCHEMA_VERSION)
probe = probe_control_plane_db(self.db_path)
self.assertEqual(probe.status, STATUS_OK)
self.assertEqual(
probe.metadata["schema_version"], control_plane_db.SCHEMA_VERSION
)
self.assertEqual(probe.metadata["active_leases"], 1)
def test_mismatched_schema_is_degraded(self):
self._build_db(control_plane_db.SCHEMA_VERSION + 99)
probe = probe_control_plane_db(self.db_path)
self.assertEqual(probe.status, STATUS_DEGRADED)
def test_probe_does_not_create_a_database(self):
"""A health check must never initialise the substrate it inspects."""
absent = str(Path(self.tmp.name) / "never-created.db")
probe_control_plane_db(absent)
self.assertFalse(Path(absent).exists())
def test_unreadable_database_is_down_not_raised(self):
Path(self.db_path).write_text("this is not a sqlite database")
probe = probe_control_plane_db(self.db_path)
self.assertEqual(probe.status, STATUS_DOWN)
class TestRedaction(unittest.TestCase):
"""AC4 — redaction. No credential-shaped text crosses the boundary."""
def test_redacts_token_assignment(self):
cleaned = redact("failed with token=ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ012345")
self.assertNotIn("ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ012345", cleaned)
self.assertIn("[redacted]", cleaned)
def test_redacts_authorization_header_text(self):
cleaned = redact("Authorization: Bearer abcdefghijklmnopqrstuvwxyz123456")
self.assertNotIn("abcdefghijklmnopqrstuvwxyz123456", cleaned)
def test_redacts_long_opaque_strings(self):
cleaned = redact("value 0123456789abcdef0123456789abcdef here")
self.assertNotIn("0123456789abcdef0123456789abcdef", cleaned)
def test_url_userinfo_and_query_are_stripped(self):
cleaned = redact_url("https://user:[email protected]/api/v1?token=xyz")
self.assertNotIn("secretpass", cleaned)
self.assertNotIn("token=xyz", cleaned)
self.assertEqual(cleaned, "https://gitea.example.com/api/v1")
def test_url_inside_free_text_is_redacted(self):
cleaned = redact("GET https://u:[email protected]/x?token=abc failed")
self.assertNotIn("u:p@", cleaned)
self.assertNotIn("token=abc", cleaned)
def test_gitea_probe_failure_detail_is_redacted(self):
boom = RuntimeError(
"connection refused for https://user:[email protected]/api/v1/version"
)
with mock.patch("webui.system_health.get_auth_header", return_value="token x"), \
mock.patch("webui.system_health.api_request", side_effect=boom):
probe = probe_gitea("gitea.example.com")
self.assertEqual(probe.status, STATUS_DOWN)
self.assertNotIn("hunter2", probe.detail)
self.assertEqual(scan_text_for_client_secrets(probe.detail), [])
def test_credential_guard_refusal_is_a_status_not_a_crash(self):
with mock.patch(
"webui.system_health.get_auth_header",
side_effect=RuntimeError("daemon guard refused"),
):
probe = probe_gitea("gitea.example.com")
self.assertEqual(probe.status, STATUS_DEGRADED)
self.assertFalse(probe.required)
class TestNamespaceSummaries(unittest.TestCase):
"""A web process cannot prove IDE namespace health, and must not claim to."""
def test_every_namespace_reports_unproven(self):
rows = namespace_summaries()
self.assertTrue(rows)
for row in rows:
with self.subTest(namespace=row["namespace"]):
self.assertEqual(row["status"], "unproven")
self.assertFalse(row["ide_namespace_proven"])
self.assertIn("client_namespace", row["reason"])
class TestSystemHealthRoutes(CleanParityMixin, unittest.TestCase):
"""The HTTP surface: versioned path, status codes, read-only guard."""
def setUp(self):
super().setUp()
clear_probe_cache()
self.addCleanup(clear_probe_cache)
self.client = TestClient(create_app())
def _patch_snapshot(self, probes, daemon_head="abc123"):
snapshot = load_system_health(probes=probes, daemon_head=daemon_head)
patcher = mock.patch(
"webui.app.load_system_health",
return_value=snapshot,
)
patcher.start()
self.addCleanup(patcher.stop)
return snapshot
def test_versioned_route_is_registered(self):
self.assertEqual(API_PATH, "/api/v1/system/health")
self._patch_snapshot(_ALL_HEALTHY)
response = self.client.get(API_PATH)
self.assertEqual(response.status_code, 200)
def test_healthy_payload_shape(self):
self._patch_snapshot(_ALL_HEALTHY)
data = self.client.get(API_PATH).json()
self.assertEqual(data["status"], STATUS_OK)
self.assertTrue(data["readiness"]["ready"])
self.assertTrue(data["readiness"]["complete"])
self.assertEqual(data["api"], API_PATH)
self.assertEqual(len(data["dependencies"]), 3)
for key in ("version", "process", "stale_runtime", "mcp_namespaces"):
self.assertIn(key, data)
self.assertIn("uptime_seconds", data["process"])
self.assertIn("mutation_safe", data["stale_runtime"])
def test_degraded_dependency_returns_503(self):
probes = (
_probe("control_plane_db", STATUS_DOWN, detail="missing", kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
)
self._patch_snapshot(probes)
response = self.client.get(API_PATH)
self.assertEqual(response.status_code, 503)
data = response.json()
self.assertFalse(data["readiness"]["ready"])
self.assertTrue(data["readiness"]["reasons"])
def test_dependency_entries_expose_status_and_latency(self):
self._patch_snapshot(_ALL_HEALTHY)
data = self.client.get(API_PATH).json()
names = {entry["name"] for entry in data["dependencies"]}
self.assertEqual(names, {"control_plane_db", "repository", "gitea"})
for entry in data["dependencies"]:
with self.subTest(dependency=entry["name"]):
self.assertIn("status", entry)
self.assertIn("required", entry)
self.assertIn("latency_ms", entry)
def test_response_body_carries_no_client_secrets(self):
self._patch_snapshot(_ALL_HEALTHY)
body = self.client.get(API_PATH).text
self.assertEqual(scan_text_for_client_secrets(body), [])
def test_deep_flag_is_forwarded(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc")
with mock.patch(
"webui.app.load_system_health", return_value=snapshot
) as loader:
self.client.get(f"{API_PATH}?deep=1")
loader.assert_called_once_with(deep=True)
def test_shallow_is_the_default(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc")
with mock.patch(
"webui.app.load_system_health", return_value=snapshot
) as loader:
self.client.get(API_PATH)
loader.assert_called_once_with(deep=False)
def test_route_rejects_mutation_methods(self):
for method in ("POST", "PUT", "PATCH", "DELETE"):
with self.subTest(method=method):
response = self.client.request(method, API_PATH)
self.assertEqual(response.status_code, 405)
self.assertEqual(response.json()["error"], "read-only-mvp")
def test_default_shallow_call_skips_the_network_probe(self):
"""The expensive probe must not run unless it was asked for."""
with mock.patch("webui.system_health.probe_gitea") as probe:
snapshot = load_system_health(deep=False)
probe.assert_not_called()
gitea = next(p for p in snapshot.dependencies if p.name == "gitea")
self.assertEqual(gitea.status, STATUS_SKIPPED)
class TestHealthRouteBackwardCompatibility(unittest.TestCase):
"""`/health` is expanded additively; MVP consumers must keep working."""
def setUp(self):
self.client = TestClient(create_app())
def test_mvp_keys_are_unchanged(self):
data = self.client.get("/health").json()
self.assertEqual(data["status"], "ok")
self.assertEqual(data["service"], "mcp-control-plane-webui")
self.assertEqual(data["mode"], "read-only-mvp")
self.assertIn("timestamp", data)
self.assertEqual(data["deployment"]["mode"], "internal-operator-console")
def test_health_points_at_the_versioned_api(self):
data = self.client.get("/health").json()
self.assertEqual(data["system_health_api"], API_PATH)
self.assertIn("uptime_seconds", data)
self.assertIn("started_at", data)
def test_health_runs_no_dependency_probe(self):
"""Liveness must stay cheap: no probe, no snapshot assembly."""
with mock.patch("webui.app.load_system_health") as loader:
response = self.client.get("/health")
self.assertEqual(response.status_code, 200)
loader.assert_not_called()
class TestSnapshotSerialisation(CleanParityMixin, unittest.TestCase):
def test_snapshot_dict_is_json_serialisable(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc123")
encoded = json.dumps(snapshot_to_dict(snapshot))
self.assertIn("readiness", encoded)
if __name__ == "__main__":
unittest.main()
+364
View File
@@ -0,0 +1,364 @@
"""Tests for the workflow-event timeline model and read API (#637).
Covers the acceptance criteria: versioned schema, adaptation of control-plane
events and Gitea handoff comments, filter by issue/PR/session, redaction of
secret-like payloads, and stable pagination.
"""
import os
import sqlite3
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient
import control_plane_db
from canonical_thread_handoff import format_cth_body
from webui import timeline
from webui.app import create_app
def _seed_db(path: str) -> None:
"""Create a control-plane DB and seed scoped work_items + events."""
# Constructing ControlPlaneDB runs the schema migration once.
control_plane_db.ControlPlaneDB(db_path=path)
conn = sqlite3.connect(path)
try:
conn.execute(
"INSERT INTO work_items(remote, org, repo, kind, number, state, updated_at) "
"VALUES (?,?,?,?,?,?,?)",
("prgs", "Scaled-Tech-Consulting", "Gitea-Tools", "issue", 637, "open", "2026-07-23T00:00:00Z"),
)
issue_wid = conn.execute("SELECT last_insert_rowid()").fetchone()[0]
conn.execute(
"INSERT INTO work_items(remote, org, repo, kind, number, state, updated_at) "
"VALUES (?,?,?,?,?,?,?)",
("prgs", "Scaled-Tech-Consulting", "Gitea-Tools", "pr", 813, "open", "2026-07-23T00:00:00Z"),
)
pr_wid = conn.execute("SELECT last_insert_rowid()").fetchone()[0]
# A work item for a different repo — must never appear in prgs/Gitea-Tools scope.
conn.execute(
"INSERT INTO work_items(remote, org, repo, kind, number, state, updated_at) "
"VALUES (?,?,?,?,?,?,?)",
("dadeschools", "Other", "Elsewhere", "issue", 1, "open", "2026-07-23T00:00:00Z"),
)
other_wid = conn.execute("SELECT last_insert_rowid()").fetchone()[0]
rows = [
(issue_wid, "allocation", "assigned author work", "2026-07-23T01:00:00Z"),
(issue_wid, "lease.renew", "token=ghs_ABCDEF1234567890abcdef lease renewed", "2026-07-23T02:00:00Z"),
(pr_wid, "pr.opened", "PR opened for review", "2026-07-23T03:00:00Z"),
(other_wid, "allocation", "off-scope event", "2026-07-23T04:00:00Z"),
]
conn.executemany(
"INSERT INTO events(work_item_id, event_type, message, created_at) VALUES (?,?,?,?)",
rows,
)
conn.commit()
finally:
conn.close()
class TestSchema(unittest.TestCase):
def test_schema_is_versioned(self):
self.assertIsInstance(timeline.TIMELINE_SCHEMA_VERSION, int)
self.assertGreaterEqual(timeline.TIMELINE_SCHEMA_VERSION, 1)
def test_event_to_dict_shape(self):
ev = timeline.WorkflowEvent(
source=timeline.SOURCE_CONTROL_PLANE,
event_type="allocation",
event_key="cp:1",
timestamp="2026-07-23T01:00:00Z",
issue_number=637,
)
d = ev.to_dict()
for key in (
"source", "event_type", "event_key", "timestamp", "actor", "role",
"issue_number", "pr_number", "session_id", "tool_name", "decision",
"message", "correlation_id", "evidence_refs", "sensitive",
):
self.assertIn(key, d)
self.assertEqual(d["evidence_refs"], [])
class TestCpAdapter(unittest.TestCase):
def test_issue_and_pr_mapping(self):
rows = [
{"event_id": 1, "event_type": "allocation", "message": "x", "created_at": "2026-07-23T01:00:00Z", "kind": "issue", "number": 637},
{"event_id": 2, "event_type": "pr.opened", "message": "y", "created_at": "2026-07-23T02:00:00Z", "kind": "pr", "number": 813},
]
events = timeline.adapt_cp_events(rows)
self.assertEqual(len(events), 2)
self.assertEqual(events[0].issue_number, 637)
self.assertIsNone(events[0].pr_number)
self.assertEqual(events[0].correlation_id, "issue#637")
self.assertIsNone(events[1].issue_number)
self.assertEqual(events[1].pr_number, 813)
def test_malformed_rows_skipped(self):
rows = [
{"event_id": None, "event_type": "x", "kind": "issue", "number": 1},
{"event_id": 5, "event_type": "", "kind": "issue", "number": 1},
{"event_id": 6, "event_type": "ok", "message": "m", "created_at": None, "kind": "issue", "number": 1},
]
events = timeline.adapt_cp_events(rows)
self.assertEqual(len(events), 1)
self.assertIsNone(events[0].timestamp)
def test_sensitive_event_flagged(self):
rows = [{"event_id": 1, "event_type": "lease.renew", "message": "m", "created_at": "2026-07-23T01:00:00Z", "kind": "issue", "number": 1}]
events = timeline.adapt_cp_events(rows)
self.assertTrue(events[0].sensitive)
class TestCthAdapter(unittest.TestCase):
def test_cth_comment_becomes_event(self):
body = format_cth_body(
cth_type="Author Handoff",
status="ready",
next_owner="reviewer",
decision="implement timeline",
proof="commit abc1234 closes #637",
next_action="review PR",
ready_to_paste_prompt="Review PR #900 as reviewer",
)
comments = [{"id": 42, "body": body, "created_at": "2026-07-23T05:00:00Z", "user": {"login": "jcwalker3"}}]
events = timeline.adapt_cth_comments(comments, kind="issue", number=637)
self.assertEqual(len(events), 1)
ev = events[0]
self.assertEqual(ev.source, timeline.SOURCE_GITEA_HANDOFF)
self.assertEqual(ev.event_type, "handoff:Author Handoff")
self.assertEqual(ev.actor, "jcwalker3")
self.assertEqual(ev.issue_number, 637)
self.assertEqual(ev.event_key, "cth:issue:637:42")
self.assertIn("#637", ev.evidence_refs)
self.assertIn("abc1234", ev.evidence_refs)
def test_non_cth_comment_ignored(self):
comments = [{"id": 1, "body": "just a normal comment", "created_at": "2026-07-23T05:00:00Z", "user": {"login": "x"}}]
self.assertEqual(timeline.adapt_cth_comments(comments, kind="issue", number=1), [])
class TestRedaction(unittest.TestCase):
def test_cp_message_redacted(self):
rows = [{"event_id": 1, "event_type": "lease", "message": "token=ghs_ABCDEF1234567890abcdef here", "created_at": "2026-07-23T01:00:00Z", "kind": "issue", "number": 1}]
events = timeline.adapt_cp_events(rows)
self.assertNotIn("ghs_ABCDEF1234567890abcdef", events[0].message or "")
def test_handoff_decision_redacted(self):
body = format_cth_body(
cth_type="Blocker",
status="blocked",
next_owner="author",
decision="password=SuperSecret123! must rotate",
proof="none",
next_action="rotate",
ready_to_paste_prompt="Rotate the credential and retry",
)
comments = [{"id": 7, "body": body, "created_at": "2026-07-23T05:00:00Z", "user": {"login": "x"}}]
events = timeline.adapt_cth_comments(comments, kind="issue", number=1)
self.assertNotIn("SuperSecret123!", events[0].decision or "")
class TestFilterSortPaginate(unittest.TestCase):
def _events(self):
return [
timeline.WorkflowEvent(source="control_plane", event_type="a", event_key="cp:3", timestamp="2026-07-23T03:00:00Z", pr_number=813),
timeline.WorkflowEvent(source="control_plane", event_type="b", event_key="cp:1", timestamp="2026-07-23T01:00:00Z", issue_number=637),
timeline.WorkflowEvent(source="control_plane", event_type="c", event_key="cp:2", timestamp="2026-07-23T02:00:00Z", issue_number=637, session_id="sess-1"),
]
def test_filter_by_issue(self):
out = timeline.filter_events(self._events(), issue=637)
self.assertEqual({e.event_key for e in out}, {"cp:1", "cp:2"})
def test_filter_by_pr(self):
out = timeline.filter_events(self._events(), pr=813)
self.assertEqual([e.event_key for e in out], ["cp:3"])
def test_filter_by_session(self):
out = timeline.filter_events(self._events(), session="sess-1")
self.assertEqual([e.event_key for e in out], ["cp:2"])
def test_stable_sort_ascending(self):
out = timeline.sort_events(self._events())
self.assertEqual([e.event_key for e in out], ["cp:1", "cp:2", "cp:3"])
def test_missing_timestamp_sorts_last(self):
evs = self._events() + [
timeline.WorkflowEvent(source="control_plane", event_type="z", event_key="cp:9", timestamp=None)
]
out = timeline.sort_events(evs)
self.assertEqual(out[-1].event_key, "cp:9")
def test_pagination_windows_and_next_offset(self):
evs = timeline.sort_events(self._events())
page1 = timeline.paginate(evs, limit=2, offset=0)
self.assertEqual(len(page1.events), 2)
self.assertEqual(page1.total, 3)
self.assertEqual(page1.next_offset, 2)
page2 = timeline.paginate(evs, limit=2, offset=2)
self.assertEqual(len(page2.events), 1)
self.assertIsNone(page2.next_offset)
def test_pagination_bounds_coerced(self):
evs = self._events()
page = timeline.paginate(evs, limit=-5, offset=-3)
self.assertGreaterEqual(page.limit, 1)
self.assertEqual(page.offset, 0)
class TestCpReader(unittest.TestCase):
def test_reads_scoped_events_only(self):
import tempfile
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
events, status = timeline.read_cp_events(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools", db_path=db
)
self.assertTrue(status.ok)
# 3 scoped events; the dadeschools/Other event is excluded.
self.assertEqual(len(events), 3)
self.assertTrue(all(e.source == "control_plane" for e in events))
# Redaction applied to the token-bearing message.
joined = " ".join(e.message or "" for e in events)
self.assertNotIn("ghs_ABCDEF1234567890abcdef", joined)
def test_missing_db_degrades(self):
events, status = timeline.read_cp_events(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
db_path="/nonexistent/path/to/cp.sqlite3",
)
self.assertEqual(events, [])
self.assertFalse(status.ok)
self.assertIsNotNone(status.reason)
class TestLoadTimeline(unittest.TestCase):
def test_handoff_not_run_without_thread_filter(self):
import tempfile
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
snap = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools", db_path=db
)
d = snap.to_dict()
handoff = [s for s in d["sources"] if s["name"] == "gitea_handoff"][0]
self.assertFalse(handoff["ok"])
self.assertIn("thread-scoped", handoff["reason"])
self.assertEqual(d["schema_version"], timeline.TIMELINE_SCHEMA_VERSION)
def test_handoff_included_via_injected_source(self):
import tempfile
body = format_cth_body(
cth_type="Author Handoff", status="ready", next_owner="reviewer",
decision="d", proof="#637", next_action="review", ready_to_paste_prompt="Review PR #1 now",
)
def source(kind, number):
return [{"id": 1, "body": body, "created_at": "2026-07-23T09:00:00Z", "user": {"login": "jcwalker3"}}]
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
snap = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
issue=637, db_path=db, comment_source=source,
)
d = snap.to_dict()
handoff = [s for s in d["sources"] if s["name"] == "gitea_handoff"][0]
self.assertTrue(handoff["ok"])
self.assertEqual(handoff["count"], 1)
# Both a CP event and the handoff event for issue 637 appear, sorted.
kinds = {e["source"] for e in d["events"]}
self.assertEqual(kinds, {"control_plane", "gitea_handoff"})
def test_failing_comment_source_degrades_only_handoff(self):
import tempfile
def boom(kind, number):
raise RuntimeError("network down")
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
snap = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
issue=637, db_path=db, comment_source=boom,
)
d = snap.to_dict()
cp = [s for s in d["sources"] if s["name"] == "control_plane"][0]
handoff = [s for s in d["sources"] if s["name"] == "gitea_handoff"][0]
self.assertTrue(cp["ok"])
self.assertFalse(handoff["ok"])
self.assertIn("network down", handoff["reason"])
class TestTimelineApi(unittest.TestCase):
def setUp(self):
self._prev_db = os.environ.get(control_plane_db.DB_PATH_ENV)
self._prev_offline = os.environ.get("WEBUI_TEST_OFFLINE")
import tempfile
self._tmpdir = tempfile.TemporaryDirectory()
self._db = os.path.join(self._tmpdir.name, "cp.sqlite3")
_seed_db(self._db)
os.environ[control_plane_db.DB_PATH_ENV] = self._db
os.environ["WEBUI_TEST_OFFLINE"] = "1"
self.client = TestClient(create_app())
def tearDown(self):
if self._prev_db is None:
os.environ.pop(control_plane_db.DB_PATH_ENV, None)
else:
os.environ[control_plane_db.DB_PATH_ENV] = self._prev_db
if self._prev_offline is None:
os.environ.pop("WEBUI_TEST_OFFLINE", None)
else:
os.environ["WEBUI_TEST_OFFLINE"] = self._prev_offline
self._tmpdir.cleanup()
def test_api_returns_timeline(self):
resp = self.client.get("/api/v1/timeline")
self.assertEqual(resp.status_code, 200)
body = resp.json()
self.assertEqual(body["schema_version"], timeline.TIMELINE_SCHEMA_VERSION)
self.assertIn("events", body)
self.assertIn("pagination", body)
self.assertGreaterEqual(body["pagination"]["total"], 1)
def test_api_filter_by_issue(self):
resp = self.client.get("/api/v1/timeline?issue=637")
self.assertEqual(resp.status_code, 200)
events = resp.json()["events"]
self.assertTrue(events)
self.assertTrue(all(e["issue_number"] == 637 for e in events))
def test_api_pagination(self):
resp = self.client.get("/api/v1/timeline?limit=1&offset=0")
self.assertEqual(resp.status_code, 200)
pg = resp.json()["pagination"]
self.assertEqual(pg["limit"], 1)
self.assertEqual(len(resp.json()["events"]), 1)
if pg["total"] > 1:
self.assertTrue(pg["has_more"])
def test_api_is_read_only(self):
resp = self.client.post("/api/v1/timeline")
self.assertIn(resp.status_code, (404, 405))
def test_api_no_secret_leak(self):
resp = self.client.get("/api/v1/timeline?issue=637")
self.assertNotIn("ghs_ABCDEF1234567890abcdef", resp.text)
if __name__ == "__main__":
unittest.main()
+131 -16
View File
@@ -45,8 +45,13 @@ 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.policy_inventory import load_policy_inventory, snapshot_to_dict as policy_snapshot_to_dict
from webui.policy_views import render_policy_page
from webui.timeline import load_timeline, snapshot_to_dict as timeline_snapshot_to_dict
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"})
@@ -70,7 +75,6 @@ async def home(_request: Request) -> HTMLResponse:
"<li><strong>Projects</strong> — registry and onboarding (#427)</li>"
"<li><strong>Prompts</strong> — canonical workflow prompt library (#428)</li>"
"<li><strong>Runtime</strong> — MCP health and stale-runtime detection (#430)</li>"
"<li><strong>Policy</strong> — workflow guardrail configuration visibility (#646)</li>"
"<li><strong>Audit</strong> — final-report paste and validator preview (#431)</li>"
"<li><strong>Worktrees</strong> — branch hygiene dashboard (#432)</li>"
"<li><strong>Leases</strong> — collision and lease visibility (#433)</li>"
@@ -81,16 +85,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)))
@@ -213,17 +244,6 @@ async def api_runtime(_request: Request) -> JSONResponse:
return JSONResponse(runtime_snapshot_to_dict(load_runtime_snapshot()))
async def policy(_request: Request) -> HTMLResponse:
snapshot = load_policy_inventory()
return HTMLResponse(
render_page(title="Policy", body_html=render_policy_page(snapshot))
)
async def api_v1_policy(_request: Request) -> JSONResponse:
return JSONResponse(policy_snapshot_to_dict(load_policy_inventory()))
async def _parse_audit_form(request: Request) -> tuple[str, str | None]:
if request.method == "GET":
return "", None
@@ -391,6 +411,101 @@ async def api_console_security_model(_request: Request) -> JSONResponse:
})
def _query_int(request: Request, key: str) -> int | None:
"""Parse an optional integer query parameter; None when absent/invalid."""
raw = request.query_params.get(key)
if raw is None or not str(raw).strip():
return None
try:
return int(str(raw).strip())
except (TypeError, ValueError):
return None
def _derive_remote(host: str) -> str:
"""Map a Gitea host to its known short remote name (control-plane scope key)."""
text = (host or "").lower()
if "prgs" in text:
return "prgs"
if "dadeschools" in text:
return "dadeschools"
return text.split(".")[0] if text else ""
def _timeline_comment_source(host: str, org: str, repo: str):
"""Build a fail-soft CTH-comment fetcher for one repo, or None when offline.
Returns a callable ``(kind, number) -> list[comment]``. Credentials or
network failures raise inside the callable so ``load_timeline`` degrades the
handoff source rather than the whole timeline. Offline test mode yields no
live source so the handoff section reports ``not run``.
"""
import os
from gitea_auth import api_fetch_page, get_auth_header, repo_api_url
offline = (os.environ.get("WEBUI_TEST_OFFLINE") or "").strip().lower() in {"1", "true", "yes"}
if offline:
return None
auth = get_auth_header(host)
if not auth:
return None
def _fetch(kind: str, number: int) -> list:
segment = "pulls" if kind == "pr" else "issues"
url = f"{repo_api_url(host, org, repo)}/{segment}/{int(number)}/comments"
comments: list = []
page = 1
while page <= 20:
raw, meta = api_fetch_page(url, auth, page=page, limit=50)
comments.extend(raw)
if bool(meta["is_final_page"]):
break
page += 1
return comments
return _fetch
async def api_v1_timeline(request: Request) -> JSONResponse:
"""Read-only workflow-event timeline (#637). Filter by issue/PR/session."""
from webui.queue_loader import _host_from_url # host normalisation helper
registry, error = _load_project_registry()
if error is not None:
return JSONResponse(error.to_dict(), status_code=500)
project = registry.projects[0] if registry.projects else None
org = request.query_params.get("org") or (project.gitea_owner if project else "")
repo = request.query_params.get("repo") or (project.repo_name if project else "")
host = _host_from_url(project.remote_host) if project else ""
remote = request.query_params.get("remote") or _derive_remote(host)
if not (remote and org and repo):
return JSONResponse(
{
"error": "timeline_scope_unresolved",
"detail": "no project in registry and no remote/org/repo query params provided",
},
status_code=400,
)
comment_source = _timeline_comment_source(host, org, repo) if (host and org and repo) else None
snapshot = load_timeline(
remote=remote,
org=org,
repo=repo,
issue=_query_int(request, "issue"),
pr=_query_int(request, "pr"),
session=(request.query_params.get("session") or None),
limit=_query_int(request, "limit"),
offset=_query_int(request, "offset"),
comment_source=comment_source,
)
return JSONResponse(timeline_snapshot_to_dict(snapshot))
async def method_not_allowed(request: Request, _exc: Exception) -> Response:
path = request.url.path
if path in _AUDIT_MUTATION_PATHS and request.method == "POST":
@@ -413,6 +528,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"]),
@@ -429,8 +545,7 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/api/prompts", api_prompts, methods=["GET"]),
Route("/runtime", runtime, methods=["GET"]),
Route("/api/runtime", api_runtime, methods=["GET"]),
Route("/policy", policy, methods=["GET"]),
Route("/api/v1/policy", api_v1_policy, methods=["GET"]),
Route("/api/v1/timeline", api_v1_timeline, methods=["GET"]),
Route("/audit", audit, methods=["GET", "POST"]),
Route("/api/audit", api_audit, methods=["GET", "POST"]),
Route("/worktrees", worktrees, methods=["GET"]),
-1
View File
@@ -8,7 +8,6 @@ NAV_ITEMS = (
("/projects", "Projects"),
("/prompts", "Prompts"),
("/runtime", "Runtime"),
("/policy", "Policy"),
("/audit", "Audit"),
("/worktrees", "Worktrees"),
("/leases", "Leases"),
-387
View File
@@ -1,387 +0,0 @@
"""Read-only workflow policy and guardrail inventory for the web UI (#646).
Policy and guardrails live in code, profiles, docs, and skills. An operator
cannot *see* the active workflow policy configuration from the console without
reading the repository tree. This module projects the major guardrails into a
redacted, machine-readable inventory with source attribution (file / module /
doc), so the console can render them as HTML tables with source pointers.
Design constraints (Phase 3, #646):
- **Read-only projection.** Nothing here edits policy or exposes a toggle that
could weaken a gate. It reports what is already enforced elsewhere.
- **Source attribution without secrets.** Every guardrail carries pointers to
the file/module/doc that owns it. Live values are compact summaries derived
from the safe policy accessors that already exist (``rbac_matrix``,
``redaction_policy``, ``audit_policy``); raw regex, tokens, and endpoints are
never embedded.
- **Redact before emit.** ``snapshot_to_dict`` runs the whole payload through
``console_redaction.redact_payload`` so a planted or accidental secret in any
projected value degrades to the placeholder rather than reaching a client.
- **Fail soft.** A projection that raises is recorded as a per-entry error and
never takes the page down; a guardrail is still listed with its sources.
- **Diff vs documented defaults where feasible.** When a guardrail declares a
documented invariant, the active projection is compared against it and the
result is reported; otherwise the diff is explicitly ``None`` with a reason.
"""
from __future__ import annotations
from dataclasses import dataclass
from typing import Any, Callable
from webui import console_audit
from webui import console_authz
from webui import console_redaction
SCHEMA_VERSION = 1
READ_ONLY_NOTE = (
"Read-only projection of guardrails enforced in code, profiles, docs, and "
"skills. This view never edits policy and exposes no gate-weakening toggle."
)
@dataclass(frozen=True)
class SourcePointer:
"""Where a guardrail is defined. Attribution only — never a secret."""
label: str
path: str
kind: str # "module" | "doc" | "script" | "config"
anchor: str | None = None
def to_dict(self) -> dict[str, Any]:
return {
"label": self.label,
"path": self.path,
"kind": self.kind,
"anchor": self.anchor,
}
@dataclass(frozen=True)
class PolicyEntry:
key: str
title: str
category: str
summary: str
sources: tuple[SourcePointer, ...]
active: dict[str, Any] | None
documented_default: dict[str, Any] | None
diff: dict[str, Any] | None
error: str | None = None
def to_dict(self) -> dict[str, Any]:
return {
"key": self.key,
"title": self.title,
"category": self.category,
"summary": self.summary,
"sources": [s.to_dict() for s in self.sources],
"active": self.active,
"documented_default": self.documented_default,
"diff": self.diff,
"error": self.error,
}
@dataclass(frozen=True)
class PolicyInventorySnapshot:
schema_version: int
read_only: bool
note: str
entries: tuple[PolicyEntry, ...]
categories: tuple[str, ...]
build_errors: tuple[str, ...]
def _diff_active_vs_default(
active: dict[str, Any] | None,
documented_default: dict[str, Any] | None,
) -> dict[str, Any] | None:
"""Compare only the keys the documented default declares.
Returns ``None`` when no documented default is declared (diff not feasible)
or when the active projection is unavailable. Otherwise reports, per
declared key, whether the active value matches the documented invariant.
"""
if not documented_default:
return None
if not active:
return {"status": "active_unavailable", "checked": {}}
checked: dict[str, Any] = {}
matches = True
for key, expected in documented_default.items():
observed = active.get(key)
ok = observed == expected
matches = matches and ok
checked[key] = {"expected": expected, "observed": observed, "matches": ok}
return {
"status": "matches_documented_default" if matches else "drift_detected",
"checked": checked,
}
# ── Live projections (compact, safe, fail-soft) ──────────────────────────────
# Each returns a small dict of already-safe machine values. They are module
# level so tests can substitute one to prove the redaction pass runs.
def _project_role_separation() -> dict[str, Any]:
matrix = console_authz.rbac_matrix()
return {
"model_version": matrix.get("model_version"),
"active_phase": matrix.get("active_phase"),
"roles": [r.get("role") for r in matrix.get("roles", [])],
"privileged_action_count": len(matrix.get("privileged_actions", [])),
"default_decision": matrix.get("default_decision"),
"execution_enabled": matrix.get("execution_enabled"),
}
def _project_redaction() -> dict[str, Any]:
policy = console_redaction.redaction_policy()
return {
"policy_version": policy.get("policy_version"),
"placeholder": policy.get("placeholder"),
"applies_to": policy.get("applies_to"),
"console_detector_count": len(policy.get("console_rules", [])),
"redact_before_persist": policy.get("redact_before_persist"),
"failure_mode": policy.get("failure_mode"),
}
def _project_audit() -> dict[str, Any]:
policy = console_audit.audit_policy()
return {
"schema_version": policy.get("schema_version"),
"required_field_count": len(policy.get("required_fields", [])),
"results": policy.get("results"),
"retention_defaults_days": policy.get("retention_defaults_days"),
"append_only": policy.get("append_only"),
"redact_before_persist": policy.get("redact_before_persist"),
"enabled": policy.get("enabled"),
}
def _static(value: dict[str, Any]) -> Callable[[], dict[str, Any]]:
return lambda: dict(value)
# ── Guardrail catalog ────────────────────────────────────────────────────────
# One row per major guardrail. ``project`` yields the active value (may raise;
# caught per entry). ``documented_default`` drives the feasible diff.
_CatalogRow = tuple[
str,
str,
str,
str,
tuple[SourcePointer, ...],
Callable[[], dict[str, Any]] | None,
dict[str, Any] | None,
]
_CATALOG: tuple[_CatalogRow, ...] = (
(
"role_separation",
"Role separation and RBAC",
"role_separation",
"Author, reviewer, merger, and reconciler capabilities are disjoint and "
"role-exclusive; self-review and self-merge are always blocked. The "
"console RBAC model defaults to deny.",
(
SourcePointer("task capability map", "task_capability_map.py", "module"),
SourcePointer("role/namespace gate", "role_namespace_gate.py", "module"),
SourcePointer("console RBAC", "webui/console_authz.py", "module"),
),
_project_role_separation,
{"default_decision": "deny", "execution_enabled": False},
),
(
"lease_rules",
"Issue and PR lease lifecycle",
"lease_rules",
"Durable work is claimed through issue locks and control-plane leases "
"with freshness, expiry, and dead-session recovery; abandoned or stale "
"claims are reclaimed only through the sanctioned recovery path.",
(
SourcePointer("issue lock store", "issue_lock_store.py", "module"),
SourcePointer("branch cleanup guard", "branch_cleanup_guard.py", "module"),
SourcePointer("safety model §5", "docs/safety-model.md", "doc", "5-mutation-gating"),
),
None,
None,
),
(
"worktree_rules",
"Author worktree binding",
"worktree_rules",
"Author mutations require a validated worktree under branches/ derived "
"from the active issue lock; silent fallback to the stable control "
"checkout or master is forbidden (#618).",
(
SourcePointer("author worktree gate", "author_mutation_worktree.py", "module"),
SourcePointer("worktree bootstrap", "scripts/worktree-start", "script"),
SourcePointer("workflow scope guard", "workflow_scope_guard.py", "module"),
),
None,
None,
),
(
"merge_confirmation",
"Explicit merge confirmation",
"merge_confirmation",
"A merge fails closed unless the caller passes the exact confirmation "
"phrase for that PR; reviewing never implies merging.",
(
SourcePointer("merge path", "merge_pr.py", "module"),
SourcePointer("merge tool gate", "gitea_mcp_server.py", "module"),
),
_static({"required_confirmation_format": "MERGE PR <n>", "auto_merge": False}),
{"auto_merge": False},
),
(
"redaction",
"Secret redaction",
"redaction",
"Every console surface runs the shared gitea_audit pass then console "
"patterns before any payload, HTML, log line, or audit record leaves "
"the server; unredactable values fail closed to the placeholder.",
(
SourcePointer("console redaction", "webui/console_redaction.py", "module"),
SourcePointer("shared redaction", "gitea_audit.py", "module"),
SourcePointer("safety model §3", "docs/safety-model.md", "doc", "3-secret-redaction"),
),
_project_redaction,
{"redact_before_persist": True},
),
(
"contamination",
"Contamination containment",
"contamination",
"A session contaminated by a direct stable-branch push or a manual MCP "
"daemon kill is blocked from review, merge, close, and completion "
"mutations until cleared (reconciler-exempt).",
(
SourcePointer("contamination gates", "gitea_mcp_server.py", "module"),
SourcePointer("stable-branch audit", "workflow_scope_guard.py", "module"),
),
None,
None,
),
(
"allocator_policy",
"Work allocation policy",
"allocator_policy",
"Workers do not self-select exclusive work; the controller-owned "
"allocator ranks the complete queue by priority then PRs-before-issues "
"then ascending number, honoring dependency edges and foreign claims.",
(
SourcePointer("allocator", "gitea_mcp_server.py", "module"),
SourcePointer("safety model §5", "docs/safety-model.md", "doc", "5-mutation-gating"),
),
_static(
{
"self_select_exclusive_work": False,
"ranking": "priority desc, PRs before issues, number asc",
"respects_dependency_edges": True,
"respects_foreign_claims": True,
}
),
{"self_select_exclusive_work": False},
),
(
"audit_logging",
"Audit logging",
"audit_logging",
"Console intent and authorization outcomes are recorded to an "
"append-only, redact-before-persist audit log; MCP mutations are "
"recorded by gitea_audit and correlated by request id.",
(
SourcePointer("console audit", "webui/console_audit.py", "module"),
SourcePointer("MCP audit", "gitea_audit.py", "module"),
SourcePointer("safety model §1", "docs/safety-model.md", "doc", "1-audit-logging-and-confirmation"),
),
_project_audit,
{"append_only": True, "redact_before_persist": True},
),
(
"mutation_gating",
"Mutation gating and master parity",
"mutation_gating",
"Mutations fail closed while the running server is stale relative to "
"master, and every mutation is preceded by identity and capability "
"resolution in a fixed pre-flight order.",
(
SourcePointer("mutation gate", "gitea_mcp_server.py", "module"),
SourcePointer("safety model §5", "docs/safety-model.md", "doc", "5-mutation-gating"),
),
_static(
{
"stale_runtime_blocks_mutations": True,
"preflight_order": "whoami -> resolve_task_capability -> mutation",
}
),
{"stale_runtime_blocks_mutations": True},
),
)
def _build_entry(row: _CatalogRow) -> PolicyEntry:
key, title, category, summary, sources, project, documented_default = row
active: dict[str, Any] | None = None
error: str | None = None
if project is not None:
try:
active = project()
except Exception as exc: # noqa: BLE001 — fail soft; never take the page down
active = None
error = f"active projection unavailable: {exc}"
diff = _diff_active_vs_default(active, documented_default)
return PolicyEntry(
key=key,
title=title,
category=category,
summary=summary,
sources=sources,
active=active,
documented_default=documented_default,
diff=diff,
error=error,
)
def load_policy_inventory() -> PolicyInventorySnapshot:
"""Build the read-only guardrail inventory. Never raises for one bad entry."""
entries: list[PolicyEntry] = []
build_errors: list[str] = []
for row in _CATALOG:
try:
entries.append(_build_entry(row))
except Exception as exc: # noqa: BLE001 — one row must not break the rest
build_errors.append(f"{row[0]}: {exc}")
categories = tuple(dict.fromkeys(e.category for e in entries))
return PolicyInventorySnapshot(
schema_version=SCHEMA_VERSION,
read_only=True,
note=READ_ONLY_NOTE,
entries=tuple(entries),
categories=categories,
build_errors=tuple(build_errors),
)
def snapshot_to_dict(snapshot: PolicyInventorySnapshot) -> dict[str, Any]:
"""Serialize the snapshot, redacting the entire payload before it is emitted."""
payload = {
"schema_version": snapshot.schema_version,
"read_only": snapshot.read_only,
"note": snapshot.note,
"categories": list(snapshot.categories),
"entry_count": len(snapshot.entries),
"entries": [entry.to_dict() for entry in snapshot.entries],
"build_errors": list(snapshot.build_errors),
}
return console_redaction.redact_payload(payload)
-104
View File
@@ -1,104 +0,0 @@
"""HTML views for the workflow policy and guardrail inventory (#646)."""
from __future__ import annotations
import html
import json
from webui.policy_inventory import PolicyEntry, PolicyInventorySnapshot
def _source_pointer(source) -> str:
path = source.path
if source.anchor:
path = f"{path}#{source.anchor}"
return (
f"<li>{html.escape(source.label)}"
f"<code>{html.escape(path)}</code> "
f"<span class='muted'>({html.escape(source.kind)})</span></li>"
)
def _active_block(entry: PolicyEntry) -> str:
if entry.error:
return (
"<p class='muted'><strong>Active value unavailable:</strong> "
f"{html.escape(entry.error)}</p>"
)
if not entry.active:
return "<p class='muted'>No live projection for this guardrail.</p>"
pretty = json.dumps(entry.active, indent=2, sort_keys=True, default=str)
return f"<pre class='prompt-text'>{html.escape(pretty)}</pre>"
def _diff_block(entry: PolicyEntry) -> str:
if entry.diff is None:
if entry.documented_default is None:
return "<p class='muted'>Diff vs documented default: not feasible (no declared default).</p>"
return "<p class='muted'>Diff vs documented default: unavailable.</p>"
status = entry.diff.get("status", "unknown")
badge = "badge-claimed" if status == "matches_documented_default" else "badge-blocked"
rows = []
for key, cell in (entry.diff.get("checked") or {}).items():
marker = "" if cell.get("matches") else ""
rows.append(
"<tr>"
f"<td><code>{html.escape(str(key))}</code></td>"
f"<td><code>{html.escape(str(cell.get('expected')))}</code></td>"
f"<td><code>{html.escape(str(cell.get('observed')))}</code></td>"
f"<td>{marker}</td>"
"</tr>"
)
table = ""
if rows:
table = (
"<table class='detail'><thead><tr>"
"<th>Key</th><th>Documented</th><th>Active</th><th>Match</th>"
"</tr></thead><tbody>"
f"{''.join(rows)}</tbody></table>"
)
return (
f"<p class='meta'>Diff vs documented default: "
f"<span class='badge {badge}'>{html.escape(status)}</span></p>"
f"{table}"
)
def _entry_card(entry: PolicyEntry) -> str:
sources = "".join(_source_pointer(s) for s in entry.sources)
return (
"<div class='prompt-card'>"
f"<h3>{html.escape(entry.title)} "
f"<span class='badge'>{html.escape(entry.category)}</span></h3>"
f"<p>{html.escape(entry.summary)}</p>"
"<p class='meta'><strong>Source pointers</strong></p>"
f"<ul>{sources}</ul>"
"<p class='meta'><strong>Active configuration</strong></p>"
f"{_active_block(entry)}"
f"{_diff_block(entry)}"
"</div>"
)
def render_policy_page(snapshot: PolicyInventorySnapshot) -> str:
categories = ", ".join(html.escape(c) for c in snapshot.categories) or "none"
cards = "".join(_entry_card(e) for e in snapshot.entries)
build_errors = ""
if snapshot.build_errors:
items = "".join(
f"<li>{html.escape(err)}</li>" for err in snapshot.build_errors
)
build_errors = (
"<div class='stub'><p><strong>Some guardrails could not be built:"
f"</strong></p><ul>{items}</ul></div>"
)
return (
"<h2>Workflow policy &amp; guardrails</h2>"
f"<p class='muted'>{html.escape(snapshot.note)}</p>"
f"<p class='meta'>Schema v{snapshot.schema_version} · "
f"{len(snapshot.entries)} guardrails · categories: {categories}</p>"
f"{build_errors}"
f"{cards}"
"<p class='muted'>This page is read-only. It reports enforced policy "
"and never edits or weakens a gate. Secret values are redacted.</p>"
)
+682
View File
@@ -0,0 +1,682 @@
"""Read-only system-health model for the operator console API (#634).
`/health` answers liveness only. Operators automating readiness checks need a
structured view of *why* the control plane is or is not usable: which
dependencies answered, how long they took, what version of the code is running,
and whether the runtime is stale relative to its remote.
Three rules shape this module.
* **Read-only.** Every probe opens its subject read-only. The control-plane
database is opened through a ``mode=ro`` URI so a health check can never
create or migrate a schema, and no probe writes, restarts, or reloads
anything — restart controls are Phase 2, and #630 forbids process-kill
recovery outright.
* **Fail-soft.** A dependency that is unreachable is a *status*, not an
exception. Probes catch their own failures and report them as a degraded or
down entry carrying a reason.
* **Never claim more than was proven.** Readiness is derived only from probes
that actually ran, ``mutation_safe`` stays false unless the parity commits are
known and equal, and an MCP namespace is reported unproven because a web
process cannot exercise the IDE-managed client path (#543).
"""
from __future__ import annotations
import os
import re
import sqlite3
import subprocess
import time
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Callable
from urllib.parse import urlsplit, urlunsplit
import control_plane_db
import mcp_namespace_health
from gitea_auth import api_request, get_auth_header, gitea_url
from webui.project_registry import load_registry
SERVICE_NAME = "mcp-control-plane-webui"
API_PATH = "/api/v1/system/health"
STATUS_OK = "ok"
STATUS_DEGRADED = "degraded"
STATUS_DOWN = "down"
STATUS_SKIPPED = "skipped"
STATUS_UNPROVEN = "unproven"
# Statuses that count as a healthy answer from a probe.
_HEALTHY_STATUSES = frozenset({STATUS_OK})
# Statuses meaning "this probe did not run", as opposed to "it ran and failed".
_NOT_RUN_STATUSES = frozenset({STATUS_SKIPPED})
_DEEP_PROBE_TTL_ENV = "WEBUI_HEALTH_PROBE_TTL_SECONDS"
_DEFAULT_DEEP_PROBE_TTL = 15.0
_GITEA_PROBE_TIMEOUT_SECONDS = 5.0
_OFFLINE_ENV = "WEBUI_TEST_OFFLINE"
# Credential-shaped material that must never reach the browser, mirroring the
# forbidden client patterns in webui/deployment_boundary.py.
_SECRET_RE = re.compile(
r"(?i)\b(token|password|passwd|secret|authorization|bearer)\b\s*[:=]?\s*\S+"
)
_LONG_OPAQUE_RE = re.compile(r"\b[A-Za-z0-9_\-]{32,}\b")
# Captured once at import so uptime measures this process, not the request.
_STARTED_AT = datetime.now(timezone.utc)
_STARTED_MONOTONIC = time.monotonic()
# TTL cache for the expensive (network) probe only.
_deep_cache: dict[str, tuple[float, "DependencyProbe"]] = {}
@dataclass(frozen=True)
class DependencyProbe:
"""One dependency check, fail-soft, with its own latency."""
name: str
kind: str
status: str
detail: str
required: bool
latency_ms: float | None = None
metadata: dict[str, Any] | None = None
@property
def healthy(self) -> bool:
return self.status in _HEALTHY_STATUSES
@property
def ran(self) -> bool:
return self.status not in _NOT_RUN_STATUSES
@dataclass(frozen=True)
class VersionInfo:
git_sha: str | None
git_describe: str | None
control_plane_schema_version: int | None
python_version: str
known: bool
@dataclass(frozen=True)
class StaleRuntime:
"""Parity between the running code, the checkout, and the remote.
``mutation_safe`` is deliberately conservative: unknown is not safe.
"""
daemon_head: str | None
checkout_head: str | None
remote_head: str | None
stale: bool
determinable: bool
mutation_safe: bool
reasons: tuple[str, ...]
@dataclass(frozen=True)
class SystemHealthSnapshot:
status: str
ready: bool
readiness_complete: bool
readiness_reasons: tuple[str, ...]
service: str
mode: str
version: VersionInfo
started_at: str
uptime_seconds: float
timestamp: str
deep_probes_requested: bool
dependencies: tuple[DependencyProbe, ...]
mcp_namespaces: tuple[dict[str, Any], ...]
stale_runtime: StaleRuntime
probe_errors: tuple[str, ...] = ()
def process_uptime() -> tuple[str, float]:
"""Process start timestamp and uptime — in-memory, safe for `/health`."""
return _STARTED_AT.isoformat(), round(time.monotonic() - _STARTED_MONOTONIC, 3)
def _offline() -> bool:
return (os.environ.get(_OFFLINE_ENV) or "").strip().lower() in {"1", "true", "yes"}
def _repo_root() -> Path:
override = (os.environ.get("WEBUI_REPO_ROOT") or "").strip()
if override:
return Path(override).resolve()
return Path(__file__).resolve().parent.parent
def _deep_probe_ttl() -> float:
raw = (os.environ.get(_DEEP_PROBE_TTL_ENV) or "").strip()
if not raw:
return _DEFAULT_DEEP_PROBE_TTL
try:
value = float(raw)
except ValueError:
return _DEFAULT_DEEP_PROBE_TTL
return value if value >= 0 else _DEFAULT_DEEP_PROBE_TTL
def redact(text: str) -> str:
"""Strip credential-shaped material from operator-visible probe text.
Probe details carry exception strings, and an exception raised by an HTTP
client can quote the request that failed. Redaction happens here, at the
boundary where those strings become part of a browser-bound payload.
"""
if not text:
return ""
cleaned = _redact_urls(text)
cleaned = _SECRET_RE.sub(lambda m: f"{m.group(1)}=[redacted]", cleaned)
return _LONG_OPAQUE_RE.sub("[redacted]", cleaned)
def _redact_urls(text: str) -> str:
return re.sub(r"https?://\S+", lambda m: redact_url(m.group(0)), text)
def redact_url(url: str) -> str:
"""Reduce a URL to scheme://host/path — no userinfo, no query, no fragment."""
try:
parts = urlsplit(url)
except ValueError:
return "[redacted-url]"
if not parts.scheme or not parts.hostname:
return "[redacted-url]"
netloc = parts.hostname
if parts.port:
netloc = f"{netloc}:{parts.port}"
return urlunsplit((parts.scheme, netloc, parts.path, "", ""))
def _git(repo: Path, *args: str) -> str | None:
try:
completed = subprocess.run(
["git", "-C", str(repo), *args],
capture_output=True,
text=True,
check=False,
timeout=10,
)
except (OSError, subprocess.SubprocessError):
return None
if completed.returncode != 0:
return None
return (completed.stdout or "").strip() or None
def _load_version(repo: Path, *, schema_version: int | None) -> VersionInfo:
import platform
git_sha = None if _offline() else _git(repo, "rev-parse", "HEAD")
describe = None if _offline() else _git(repo, "describe", "--tags", "--always")
return VersionInfo(
git_sha=git_sha,
git_describe=describe,
control_plane_schema_version=schema_version,
python_version=platform.python_version(),
known=bool(git_sha),
)
# ---------------------------------------------------------------------------
# Dependency probes
# ---------------------------------------------------------------------------
def _elapsed_ms(started: float) -> float:
return round((time.monotonic() - started) * 1000, 3)
def probe_control_plane_db(db_path: str | None = None) -> DependencyProbe:
"""Read-only reachability check for the control-plane SQLite substrate.
Opened through a ``mode=ro`` URI on purpose: ``ControlPlaneDB.__init__``
creates directories and runs schema migrations, which a health check must
never do.
"""
path = (db_path or control_plane_db.default_db_path()).strip()
started = time.monotonic()
metadata: dict[str, Any] = {"path": path}
def _result(status: str, detail: str) -> DependencyProbe:
return DependencyProbe(
name="control_plane_db",
kind="sqlite",
status=status,
detail=detail,
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
if not path or not os.path.exists(path):
return _result(STATUS_DOWN, "control-plane database file does not exist yet")
try:
conn = sqlite3.connect(f"file:{path}?mode=ro", uri=True, timeout=5)
try:
row = conn.execute(
"SELECT value FROM schema_meta WHERE key = 'schema_version'"
).fetchone()
leases = conn.execute(
"SELECT COUNT(*) FROM leases WHERE status = 'active'"
).fetchone()
finally:
conn.close()
except sqlite3.Error as exc:
return _result(STATUS_DOWN, redact(f"control-plane database unreadable: {exc}"))
schema_version = int(row[0]) if row and str(row[0]).isdigit() else None
metadata["schema_version"] = schema_version
metadata["active_leases"] = int(leases[0]) if leases else None
if schema_version is None:
return _result(
STATUS_DEGRADED, "control-plane database has no recorded schema version"
)
if schema_version != control_plane_db.SCHEMA_VERSION:
return _result(
STATUS_DEGRADED,
f"control-plane schema version {schema_version} does not match the "
f"version this code expects ({control_plane_db.SCHEMA_VERSION})",
)
return _result(STATUS_OK, f"schema v{schema_version} readable")
def probe_repository(repo: Path) -> DependencyProbe:
"""Local checkout reachability — required, cheap, no network."""
started = time.monotonic()
metadata: dict[str, Any] = {"repo_root": str(repo)}
if _offline():
return DependencyProbe(
name="repository",
kind="git",
status=STATUS_SKIPPED,
detail=f"{_OFFLINE_ENV} is set; git probe skipped",
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
head = _git(repo, "rev-parse", "HEAD")
if not head:
return DependencyProbe(
name="repository",
kind="git",
status=STATUS_DOWN,
detail=f"HEAD could not be read at {repo}",
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
branch = _git(repo, "rev-parse", "--abbrev-ref", "HEAD")
metadata["head"] = head
metadata["branch"] = branch
return DependencyProbe(
name="repository",
kind="git",
status=STATUS_OK,
detail=f"checkout readable at {branch or 'detached HEAD'}",
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
def probe_gitea(host: str) -> DependencyProbe:
"""Live Gitea reachability. Expensive (network), so opt-in via ``deep``.
Optional by design: the console stays useful for local inventory when the
remote is unreachable, so a failure here degrades status without claiming
the process itself is unready.
"""
started = time.monotonic()
metadata: dict[str, Any] = {"host": host}
def _failure(status: str, detail: str) -> DependencyProbe:
return DependencyProbe(
name="gitea",
kind="http",
status=status,
detail=detail,
required=False,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
if not host:
return _failure(STATUS_DEGRADED, "no Gitea host is configured in the registry")
try:
auth = get_auth_header(host)
except Exception as exc: # noqa: BLE001 — credential guards are a status here
return _failure(STATUS_DEGRADED, redact(f"credential lookup refused: {exc}"))
if not auth:
return _failure(STATUS_DEGRADED, f"no credentials available for {host}")
url = gitea_url(host, "/api/v1/version")
metadata["endpoint"] = redact_url(url)
try:
data = api_request("GET", url, auth, timeout=_GITEA_PROBE_TIMEOUT_SECONDS)
except Exception as exc: # noqa: BLE001 — a down dependency is a status
return _failure(STATUS_DOWN, redact(f"Gitea probe failed: {exc}"))
if isinstance(data, dict) and data.get("version"):
metadata["gitea_version"] = str(data["version"])
return DependencyProbe(
name="gitea",
kind="http",
status=STATUS_OK,
detail=f"{host} reachable",
required=False,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
def _skipped_gitea(host: str) -> DependencyProbe:
return DependencyProbe(
name="gitea",
kind="http",
status=STATUS_SKIPPED,
detail="network probe not requested; call with ?deep=1 to run it",
required=False,
latency_ms=None,
metadata={"host": host},
)
def namespace_summaries() -> tuple[dict[str, Any], ...]:
"""Declared MCP namespaces, each honestly reported as unproven.
The web process runs outside the IDE-managed MCP client, so it cannot
invoke a namespace tool. Per #543 only a ``client_namespace`` probe proves
that path, and inventing a healthy verdict here is exactly the false claim
the mutation gates exist to prevent.
"""
rows: list[dict[str, Any]] = []
for namespace, required_tool in sorted(
mcp_namespace_health.REQUIRED_NAMESPACE_TOOLS.items()
):
classification = mcp_namespace_health.classify_namespace_probe(
namespace,
required_tool=required_tool,
probe_result=None,
probe_source=mcp_namespace_health.PROBE_SOURCE_UNKNOWN,
)
rows.append(
{
"namespace": namespace,
"required_tool": required_tool,
"status": STATUS_UNPROVEN,
"ide_namespace_proven": bool(classification.get("ide_namespace_proven")),
"reason": (
"the web console cannot invoke the IDE-managed MCP client; "
"namespace health must be proven with a client_namespace "
"probe (#543)"
),
"error_type": classification.get("error_type"),
}
)
return tuple(rows)
def assess_stale_runtime(
repo: Path,
*,
daemon_head: str | None = None,
git_reader: Callable[..., str | None] | None = None,
) -> StaleRuntime:
"""Three-way parity view: running code, local checkout, remote-tracking ref.
``mutation_safe`` requires all three to be known and equal. Anything less —
including "the remote ref was never fetched" — is reported as not safe with
a reason, so an operator never reads an unproven green.
"""
reader = git_reader or (lambda *args: _git(repo, *args))
reasons: list[str] = []
# The offline switch suppresses real subprocess calls; an explicitly
# injected reader is already a substitute for them and is always used.
offline = _offline() and git_reader is None
checkout_head = None if offline else reader("rev-parse", "HEAD")
remote_head = None if offline else reader("rev-parse", "@{upstream}")
if offline:
reasons.append(f"{_OFFLINE_ENV} is set; parity commits were not read")
else:
if checkout_head is None:
reasons.append("local checkout HEAD could not be read")
if remote_head is None:
reasons.append(
"no remote-tracking commit is known for the current branch; "
"remote staleness is indeterminate (no fetch is performed here)"
)
effective_daemon = daemon_head if daemon_head is not None else checkout_head
if daemon_head is None:
reasons.append(
"the running MCP daemon's startup commit is not observable from the "
"web process; the checkout commit is reported in its place"
)
determinable = bool(checkout_head and remote_head and effective_daemon)
stale = bool(
determinable and len({checkout_head, remote_head, effective_daemon}) > 1
)
if stale:
reasons.append(
"runtime, checkout, and remote commits disagree; restart the MCP "
"server after updating the checkout before trusting capability gates"
)
return StaleRuntime(
daemon_head=effective_daemon,
checkout_head=checkout_head,
remote_head=remote_head,
stale=stale,
determinable=determinable,
mutation_safe=bool(determinable and not stale),
reasons=tuple(reasons),
)
# ---------------------------------------------------------------------------
# Snapshot assembly
# ---------------------------------------------------------------------------
def _default_host() -> str:
registry = load_registry()
if not registry.projects:
return ""
raw = registry.projects[0].remote_host
parts = urlsplit(raw.strip())
return parts.netloc or raw.strip().rstrip("/")
def _aggregate(
probes: tuple[DependencyProbe, ...],
) -> tuple[str, bool, bool, tuple[str, ...]]:
"""Fold probe results into overall status and readiness.
Required probes drive readiness; optional probes can only degrade status.
A probe that did not run leaves readiness incomplete rather than passing.
"""
reasons: list[str] = []
required = [probe for probe in probes if probe.required]
unrun_required = [probe for probe in required if not probe.ran]
failed_required = [probe for probe in required if probe.ran and not probe.healthy]
failed_optional = [
probe
for probe in probes
if not probe.required and probe.ran and not probe.healthy
]
for probe in unrun_required:
reasons.append(
f"required dependency '{probe.name}' was not probed: {probe.detail}"
)
for probe in failed_required:
reasons.append(
f"required dependency '{probe.name}' is {probe.status}: {probe.detail}"
)
for probe in failed_optional:
reasons.append(
f"optional dependency '{probe.name}' is {probe.status}: {probe.detail}"
)
readiness_complete = not unrun_required
ready = readiness_complete and not failed_required
if any(probe.status == STATUS_DOWN for probe in failed_required):
status = STATUS_DOWN
elif failed_required or failed_optional or unrun_required:
status = STATUS_DEGRADED
else:
status = STATUS_OK
return status, ready, readiness_complete, tuple(reasons)
def load_system_health(
*,
deep: bool = False,
host: str | None = None,
probes: tuple[DependencyProbe, ...] | None = None,
daemon_head: str | None = None,
use_cache: bool = True,
) -> SystemHealthSnapshot:
"""Assemble the read-only system-health snapshot.
``deep=True`` adds the network probe against Gitea; its result is cached for
a short TTL so repeated dashboard polls do not amplify into remote load.
"""
repo = _repo_root()
probe_errors: list[str] = []
if probes is None:
collected: list[DependencyProbe] = []
for probe_fn in (
lambda: probe_control_plane_db(),
lambda: probe_repository(repo),
):
try:
collected.append(probe_fn())
except Exception as exc: # noqa: BLE001 — a probe must not 500 the API
probe_errors.append(redact(f"probe raised: {exc}"))
resolved_host = host if host is not None else _default_host()
if deep and not _offline():
collected.append(_cached_gitea_probe(resolved_host, use_cache=use_cache))
else:
collected.append(_skipped_gitea(resolved_host))
probes = tuple(collected)
status, ready, readiness_complete, reasons = _aggregate(probes)
stale = assess_stale_runtime(repo, daemon_head=daemon_head)
if stale.stale:
if status == STATUS_OK:
status = STATUS_DEGRADED
reasons = reasons + (
"runtime is stale relative to its remote-tracking commit",
)
db_probe = next((p for p in probes if p.name == "control_plane_db"), None)
schema_version = None
if db_probe and db_probe.metadata:
schema_version = db_probe.metadata.get("schema_version")
return SystemHealthSnapshot(
status=status,
ready=ready,
readiness_complete=readiness_complete,
readiness_reasons=reasons,
service=SERVICE_NAME,
mode="read-only",
version=_load_version(repo, schema_version=schema_version),
started_at=_STARTED_AT.isoformat(),
uptime_seconds=round(time.monotonic() - _STARTED_MONOTONIC, 3),
timestamp=datetime.now(timezone.utc).isoformat(),
deep_probes_requested=deep,
dependencies=probes,
mcp_namespaces=namespace_summaries(),
stale_runtime=stale,
probe_errors=tuple(probe_errors),
)
def _cached_gitea_probe(host: str, *, use_cache: bool = True) -> DependencyProbe:
ttl = _deep_probe_ttl()
now = time.monotonic()
if use_cache and ttl > 0:
cached = _deep_cache.get(host)
if cached and (now - cached[0]) < ttl:
return cached[1]
probe = probe_gitea(host)
if use_cache and ttl > 0:
_deep_cache[host] = (now, probe)
return probe
def clear_probe_cache() -> None:
"""Drop cached deep-probe results (tests and operator-forced refresh)."""
_deep_cache.clear()
def probe_to_dict(probe: DependencyProbe) -> dict[str, Any]:
return {
"name": probe.name,
"kind": probe.kind,
"status": probe.status,
"detail": probe.detail,
"required": probe.required,
"healthy": probe.healthy,
"latency_ms": probe.latency_ms,
"metadata": dict(probe.metadata or {}),
}
def snapshot_to_dict(snapshot: SystemHealthSnapshot) -> dict[str, Any]:
return {
"status": snapshot.status,
"service": snapshot.service,
"mode": snapshot.mode,
"api": API_PATH,
"timestamp": snapshot.timestamp,
"readiness": {
"ready": snapshot.ready,
"complete": snapshot.readiness_complete,
"reasons": list(snapshot.readiness_reasons),
},
"version": {
"git_sha": snapshot.version.git_sha,
"git_describe": snapshot.version.git_describe,
"control_plane_schema_version": (
snapshot.version.control_plane_schema_version
),
"python_version": snapshot.version.python_version,
"known": snapshot.version.known,
},
"process": {
"started_at": snapshot.started_at,
"uptime_seconds": snapshot.uptime_seconds,
},
"deep_probes_requested": snapshot.deep_probes_requested,
"dependencies": [probe_to_dict(probe) for probe in snapshot.dependencies],
"mcp_namespaces": [dict(row) for row in snapshot.mcp_namespaces],
"stale_runtime": {
"daemon_head": snapshot.stale_runtime.daemon_head,
"checkout_head": snapshot.stale_runtime.checkout_head,
"remote_head": snapshot.stale_runtime.remote_head,
"stale": snapshot.stale_runtime.stale,
"determinable": snapshot.stale_runtime.determinable,
"mutation_safe": snapshot.stale_runtime.mutation_safe,
"reasons": list(snapshot.stale_runtime.reasons),
},
"probe_errors": list(snapshot.probe_errors),
}
+539
View File
@@ -0,0 +1,539 @@
"""Workflow-event and conversation timeline model (#637, Phase 1).
Operators cannot browse a unified timeline of workflow events, decisions,
tool calls, and handoffs: the evidence is scattered across control-plane
events, Gitea canonical handoff comments, and local logs. This module defines
one durable, versioned event schema and per-source adapters that normalise
those scattered records into a single ``WorkflowEvent`` stream, plus a
read-only query layer (filter by issue / PR / session, stable ordering,
pagination) that the ``/api/v1/timeline`` route serves.
Design rules honoured here:
- **Read-only.** Sources are read; nothing is mutated. The control-plane
database is opened through a ``mode=ro`` URI so a missing or unwritable DB
degrades to a reason instead of creating directories or running migrations.
- **Fail-soft per source.** An unavailable source degrades to a status with a
reason rather than raising, and a source that could not run is never
rendered as an empty-and-healthy timeline.
- **Redaction at the boundary, fail closed.** Every free-text field (event
messages, redacted tool arguments, decision/proof text) is run through the
console redaction policy before it leaves this module. An unredactable value
becomes the placeholder — an unredacted payload is never emitted, and a
generation error never drops raw data to a caller or a log.
- **Stable ordering.** Events sort by ``(timestamp, source_rank, event_key)``
with a deterministic tiebreak, so pagination is stable across calls and
events with equal or missing timestamps keep a fixed order.
Non-goals (from the issue): no full chat replay, no mutation of historical
events, no unredacted tool-argument storage.
"""
from __future__ import annotations
import re
import sqlite3
from dataclasses import dataclass
from datetime import datetime, timezone
from typing import Any, Callable, Iterable
import control_plane_db
from webui import console_redaction
# The schema is versioned so consumers can branch on shape. Bump on any
# breaking change to WorkflowEvent's serialized form.
TIMELINE_SCHEMA_VERSION = 1
# Known event sources and their deterministic ordering rank. When two events
# carry the same timestamp, the source rank breaks the tie before the
# per-source event key, so a control-plane event and a handoff comment minted
# in the same second always sort in a fixed order.
SOURCE_CONTROL_PLANE = "control_plane"
SOURCE_GITEA_HANDOFF = "gitea_handoff"
_SOURCE_RANK = {
SOURCE_CONTROL_PLANE: 0,
SOURCE_GITEA_HANDOFF: 1,
}
# A timestamp far in the future so events with no parseable timestamp sort
# last (after everything real) instead of first, without raising.
_MISSING_TS_SORT = "9999-12-31T23:59:59Z"
def _parse_ts(value: str | None) -> str | None:
"""Normalise a timestamp to ``...Z`` UTC, or None when unparseable."""
if not value:
return None
text = str(value).strip()
if not text:
return None
candidate = text[:-1] + "+00:00" if text.endswith("Z") else text
try:
parsed = datetime.fromisoformat(candidate)
except ValueError:
return None
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z")
def _redact(value: Any) -> Any:
"""Redact a single free-text field, failing closed to the placeholder."""
if value is None:
return None
return console_redaction.redact_text(str(value))
@dataclass(frozen=True)
class WorkflowEvent:
"""One normalised timeline event.
Every field is optional except ``source``/``event_type``/``event_key``
because sources carry different subsets. The class is frozen so an adapted
event is an immutable record; a consumer that needs a variant builds a new
one rather than mutating history.
"""
source: str
event_type: str
event_key: str
timestamp: str | None = None
actor: str | None = None
role: str | None = None
issue_number: int | None = None
pr_number: int | None = None
session_id: str | None = None
tool_name: str | None = None
decision: str | None = None
message: str | None = None
correlation_id: str | None = None
evidence_refs: tuple[str, ...] = ()
sensitive: bool = False
def sort_key(self) -> tuple[str, int, str]:
return (
self.timestamp or _MISSING_TS_SORT,
_SOURCE_RANK.get(self.source, 99),
self.event_key,
)
def to_dict(self) -> dict[str, Any]:
return {
"source": self.source,
"event_type": self.event_type,
"event_key": self.event_key,
"timestamp": self.timestamp,
"actor": self.actor,
"role": self.role,
"issue_number": self.issue_number,
"pr_number": self.pr_number,
"session_id": self.session_id,
"tool_name": self.tool_name,
"decision": self.decision,
"message": self.message,
"correlation_id": self.correlation_id,
"evidence_refs": list(self.evidence_refs),
"sensitive": self.sensitive,
}
# --------------------------------------------------------------------------- #
# Adapters — pure functions from a source's raw records to WorkflowEvents. #
# Each is total: a malformed record is skipped, never raised on. #
# --------------------------------------------------------------------------- #
# Event types whose payload is treated as sensitive and always redaction-hard
# (they can carry lease/session provenance or tool arguments).
_SENSITIVE_EVENT_HINTS = ("lease", "capability", "token", "auth", "secret")
# Reference tokens (issue/PR/comment ids) and SHAs parsed out of proof text.
_EVIDENCE_REF_RE = re.compile(r"(?:#|PR\s*#?|issue\s*#?|comment\s*#?)(\d+)", re.IGNORECASE)
_SHA_RE = re.compile(r"\b[0-9a-f]{7,40}\b")
def _kind_to_numbers(kind: str | None, number: int | None) -> tuple[int | None, int | None]:
"""Map a control-plane work-item (kind, number) to (issue_no, pr_no)."""
if number is None:
return (None, None)
if kind == "pr":
return (None, int(number))
if kind == "issue":
return (int(number), None)
return (None, None)
def _correlation_for(kind: str | None, number: int | None) -> str | None:
if number is None or kind not in ("issue", "pr"):
return None
return f"{kind}#{number}"
def _extract_evidence_refs(*texts: str | None) -> tuple[str, ...]:
refs: list[str] = []
for text in texts:
if not text:
continue
for match in _EVIDENCE_REF_RE.finditer(text):
token = f"#{match.group(1)}"
if token not in refs:
refs.append(token)
for match in _SHA_RE.finditer(text):
token = match.group(0)
if token not in refs:
refs.append(token)
return tuple(refs)
def adapt_cp_events(rows: Iterable[dict[str, Any]]) -> list[WorkflowEvent]:
"""Adapt control-plane ``events`` rows (joined to work_items) into events.
Each row is expected to carry ``event_id``, ``event_type``, ``message``,
``created_at`` and the joined work-item ``kind``/``number``. Rows missing
an id or type are skipped so a partially written table never raises.
"""
events: list[WorkflowEvent] = []
for row in rows or []:
try:
event_id = row.get("event_id")
event_type = (row.get("event_type") or "").strip()
if event_id is None or not event_type:
continue
kind = row.get("kind")
number = row.get("number")
issue_no, pr_no = _kind_to_numbers(kind, number)
sensitive = any(hint in event_type.lower() for hint in _SENSITIVE_EVENT_HINTS)
events.append(
WorkflowEvent(
source=SOURCE_CONTROL_PLANE,
event_type=event_type,
event_key=f"cp:{event_id}",
timestamp=_parse_ts(row.get("created_at")),
issue_number=issue_no,
pr_number=pr_no,
session_id=(row.get("session_id") or None),
message=_redact(row.get("message")),
correlation_id=_correlation_for(kind, number),
sensitive=sensitive,
)
)
except Exception:
# A single malformed row must not sink the whole adaptation.
continue
return events
def adapt_cth_comments(
comments: Iterable[dict[str, Any]],
*,
kind: str,
number: int,
) -> list[WorkflowEvent]:
"""Adapt Gitea Canonical Thread Handoff (CTH) comments into events.
Only comments that parse as a CTH (``canonical_thread_handoff.parse_cth_comment``)
become events; ordinary comments are ignored. ``kind``/``number`` scope the
events to the issue or PR the comments belong to.
"""
# Imported lazily so this module has no import-time dependency on the
# handoff parser when only the control-plane adapter is used.
from canonical_thread_handoff import parse_cth_comment
issue_no, pr_no = _kind_to_numbers(kind, number)
correlation = _correlation_for(kind, number)
events: list[WorkflowEvent] = []
for comment in comments or []:
try:
body = comment.get("body") or ""
parsed = parse_cth_comment(body)
if not parsed:
continue
fields = parsed.get("fields") or {}
cth_type = parsed.get("cth_type") or "handoff"
comment_id = comment.get("id")
actor = (comment.get("user") or {}).get("login")
decision = fields.get("decision")
proof = fields.get("proof")
next_action = fields.get("next action")
events.append(
WorkflowEvent(
source=SOURCE_GITEA_HANDOFF,
event_type=f"handoff:{cth_type}",
event_key=f"cth:{kind}:{number}:{comment_id}",
timestamp=_parse_ts(comment.get("created_at")),
actor=actor,
role=_redact(fields.get("next owner")),
issue_number=issue_no,
pr_number=pr_no,
decision=_redact(decision),
message=_redact(next_action or fields.get("status")),
correlation_id=correlation,
evidence_refs=_extract_evidence_refs(proof, decision),
sensitive=False,
)
)
except Exception:
continue
return events
# --------------------------------------------------------------------------- #
# Read-only control-plane event source. #
# --------------------------------------------------------------------------- #
_CP_EVENTS_QUERY = """
SELECT e.event_id AS event_id,
e.event_type AS event_type,
e.message AS message,
e.created_at AS created_at,
w.kind AS kind,
w.number AS number
FROM events e
JOIN work_items w ON e.work_item_id = w.work_item_id
WHERE w.remote = ? AND w.org = ? AND w.repo = ?
"""
@dataclass(frozen=True)
class SourceStatus:
"""Fail-soft status for one timeline source."""
name: str
ok: bool
reason: str | None = None
count: int = 0
def to_dict(self) -> dict[str, Any]:
return {"name": self.name, "ok": self.ok, "reason": self.reason, "count": self.count}
def read_cp_events(
*,
remote: str,
org: str,
repo: str,
db_path: str | None = None,
) -> tuple[list[WorkflowEvent], SourceStatus]:
"""Read scoped control-plane events read-only. Never creates the DB.
Opens the SQLite file through a ``mode=ro`` URI: a health/timeline read
must never create directories or run the schema migration that
``ControlPlaneDB()`` performs on construction. A missing or unreadable DB
degrades to a status with a reason.
"""
path = (db_path or control_plane_db.default_db_path()).strip()
conn: sqlite3.Connection | None = None
try:
conn = sqlite3.connect(f"file:{path}?mode=ro", uri=True)
conn.row_factory = sqlite3.Row
cursor = conn.execute(_CP_EVENTS_QUERY, (remote, org, repo))
rows = [dict(r) for r in cursor.fetchall()]
except sqlite3.OperationalError as exc:
return ([], SourceStatus(SOURCE_CONTROL_PLANE, ok=False, reason=f"control-plane DB unavailable: {exc}"))
except sqlite3.Error as exc:
return ([], SourceStatus(SOURCE_CONTROL_PLANE, ok=False, reason=f"control-plane read failed: {exc}"))
finally:
if conn is not None:
conn.close()
events = adapt_cp_events(rows)
return (events, SourceStatus(SOURCE_CONTROL_PLANE, ok=True, count=len(events)))
# --------------------------------------------------------------------------- #
# Filter, sort, paginate. #
# --------------------------------------------------------------------------- #
def filter_events(
events: Iterable[WorkflowEvent],
*,
issue: int | None = None,
pr: int | None = None,
session: str | None = None,
) -> list[WorkflowEvent]:
"""Filter events by issue number, PR number, and/or session id.
Filters are conjunctive. A filter that names a dimension an event does not
carry excludes that event (an issue filter excludes PR-only events).
"""
out: list[WorkflowEvent] = []
for ev in events:
if issue is not None and ev.issue_number != issue:
continue
if pr is not None and ev.pr_number != pr:
continue
if session is not None and ev.session_id != session:
continue
out.append(ev)
return out
def sort_events(events: Iterable[WorkflowEvent]) -> list[WorkflowEvent]:
"""Return events in stable timeline order (ascending)."""
return sorted(events, key=lambda ev: ev.sort_key())
@dataclass(frozen=True)
class TimelinePage:
"""One page of the sorted, filtered timeline."""
events: tuple[WorkflowEvent, ...]
total: int
limit: int
offset: int
@property
def next_offset(self) -> int | None:
nxt = self.offset + len(self.events)
return nxt if nxt < self.total else None
def to_dict(self) -> dict[str, Any]:
return {
"events": [ev.to_dict() for ev in self.events],
"pagination": {
"total": self.total,
"limit": self.limit,
"offset": self.offset,
"returned": len(self.events),
"next_offset": self.next_offset,
"has_more": self.next_offset is not None,
},
}
_MAX_LIMIT = 500
_DEFAULT_LIMIT = 50
def _coerce_bounds(limit: int | None, offset: int | None) -> tuple[int, int]:
try:
lim = int(limit) if limit is not None else _DEFAULT_LIMIT
except (TypeError, ValueError):
lim = _DEFAULT_LIMIT
try:
off = int(offset) if offset is not None else 0
except (TypeError, ValueError):
off = 0
lim = max(1, min(lim, _MAX_LIMIT))
off = max(0, off)
return (lim, off)
def paginate(events: list[WorkflowEvent], *, limit: int | None, offset: int | None) -> TimelinePage:
lim, off = _coerce_bounds(limit, offset)
window = events[off : off + lim]
return TimelinePage(events=tuple(window), total=len(events), limit=lim, offset=off)
# --------------------------------------------------------------------------- #
# Composition — load_timeline aggregates all sources, fail-soft. #
# --------------------------------------------------------------------------- #
# A comment source is a callable that, given (kind, number), returns the raw
# Gitea comment list for that issue/PR. The route supplies a live fail-soft
# fetcher; tests supply a fixture. When None, the handoff source is reported as
# not-run (never silently empty-and-healthy).
CommentSource = Callable[[str, int], list[dict[str, Any]]]
@dataclass(frozen=True)
class TimelineSnapshot:
schema_version: int
remote: str
org: str
repo: str
filters: dict[str, Any]
page: TimelinePage
sources: tuple[SourceStatus, ...]
def to_dict(self) -> dict[str, Any]:
return {
"schema_version": self.schema_version,
"scope": {"remote": self.remote, "org": self.org, "repo": self.repo},
"filters": self.filters,
"sources": [s.to_dict() for s in self.sources],
**self.page.to_dict(),
}
def load_timeline(
*,
remote: str,
org: str,
repo: str,
issue: int | None = None,
pr: int | None = None,
session: str | None = None,
limit: int | None = None,
offset: int | None = None,
db_path: str | None = None,
comment_source: CommentSource | None = None,
) -> TimelineSnapshot:
"""Aggregate every timeline source into one filtered, paginated snapshot.
Sources are read independently and fail soft: an unavailable source
contributes a ``SourceStatus`` with ``ok=False`` and a reason, and never
collapses the whole timeline. The handoff source only runs when a specific
issue or PR is requested (a handoff comment belongs to one thread) and a
``comment_source`` is available; otherwise it is reported as ``not run``
rather than as an empty-and-healthy source.
"""
all_events: list[WorkflowEvent] = []
statuses: list[SourceStatus] = []
cp_events, cp_status = read_cp_events(remote=remote, org=org, repo=repo, db_path=db_path)
all_events.extend(cp_events)
statuses.append(cp_status)
# Gitea handoff comments are thread-scoped: only fetch when the caller
# narrowed to one issue or PR, and only when a source was provided.
handoff_target: tuple[str, int] | None = None
if pr is not None:
handoff_target = ("pr", pr)
elif issue is not None:
handoff_target = ("issue", issue)
if handoff_target is None:
statuses.append(
SourceStatus(
SOURCE_GITEA_HANDOFF,
ok=False,
reason="not run: handoff comments are thread-scoped; filter by issue or pr to include them",
)
)
elif comment_source is None:
statuses.append(
SourceStatus(
SOURCE_GITEA_HANDOFF,
ok=False,
reason="not run: no comment source configured for this timeline read",
)
)
else:
kind, number = handoff_target
try:
comments = comment_source(kind, number) or []
handoff_events = adapt_cth_comments(comments, kind=kind, number=number)
all_events.extend(handoff_events)
statuses.append(SourceStatus(SOURCE_GITEA_HANDOFF, ok=True, count=len(handoff_events)))
except Exception as exc: # fail soft: a fetch/parse error degrades this source only
statuses.append(
SourceStatus(SOURCE_GITEA_HANDOFF, ok=False, reason=f"handoff source failed: {exc}")
)
filtered = filter_events(all_events, issue=issue, pr=pr, session=session)
ordered = sort_events(filtered)
page = paginate(ordered, limit=limit, offset=offset)
return TimelineSnapshot(
schema_version=TIMELINE_SCHEMA_VERSION,
remote=remote,
org=org,
repo=repo,
filters={"issue": issue, "pr": pr, "session": session},
page=page,
sources=tuple(statuses),
)
def snapshot_to_dict(snapshot: TimelineSnapshot) -> dict[str, Any]:
return snapshot.to_dict()