Compare commits

..
Author SHA1 Message Date
sysadmin 996e7094fe chore: merge master into feat/issue-646-policy-guardrail-visibility (base sync) 2026-07-23 17:42:49 -04:00
jcwalker3andClaude Opus 4.8 ab33337a94 feat(webui): read-only workflow policy & guardrail visibility (Closes #646)
Phase 3 child of the Web Console epic #631. Operators can now see the active
workflow policy/guardrail configuration from the console instead of reading the
repo tree.

- webui/policy_inventory.py (new): redacted, machine-readable guardrail
  inventory. One row per major guardrail (role separation/RBAC, lease rules,
  worktree binding, merge confirmation, redaction, contamination, allocator
  policy, audit logging, mutation gating) with source pointers (file/module/doc)
  and a compact active projection from the existing safe policy accessors.
  Fail-soft per entry; whole payload run through console_redaction before emit;
  diff vs documented default where feasible.
- webui/policy_views.py (new): HTML cards with source pointers, active config,
  and the documented-default diff; read-only page copy, no forms.
- webui/app.py: register GET /policy and GET /api/v1/policy (additive).
- webui/layout.py: add Policy nav item.
- tests/test_webui_policy_visibility.py (new): guardrail presence + source
  pointers (AC1), redaction incl. planted-secret masking and scan_for_secrets
  (AC2/AC3), read-only page + no-mutation routes (AC4), fail-soft rendering.
- docs/webui-local-dev.md: route table + read-only policy-visibility section.

Read-only throughout; no policy editing, no gate-weakening toggle, secrets
redacted.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-23 03:20:56 -05:00
15 changed files with 872 additions and 3184 deletions
+2 -19
View File
@@ -46,17 +46,6 @@ _FIELD_RE = re.compile(
) )
def is_known_cth_type(value: str | None) -> bool:
"""True when *value* is a declared member of the :data:`CTH_TYPES` contract.
``CTH_TYPES`` is the single authority for what a CTH type may be. The
heading a comment carries is free text, so a *read* path that turns a parsed
type into something durable — a serialized field, a routing decision — must
check membership here rather than trust the parse or keep a list of its own.
"""
return (value or "").strip() in CTH_TYPES
def format_cth_body( def format_cth_body(
*, *,
cth_type: str, cth_type: str,
@@ -71,7 +60,7 @@ def format_cth_body(
) -> str: ) -> str:
"""Render a canonical CTH comment body.""" """Render a canonical CTH comment body."""
normalized_type = (cth_type or "").strip() normalized_type = (cth_type or "").strip()
if not is_known_cth_type(normalized_type): if normalized_type not in CTH_TYPES:
raise ValueError( raise ValueError(
f"unknown CTH type '{cth_type}'; expected one of {sorted(CTH_TYPES)}" f"unknown CTH type '{cth_type}'; expected one of {sorted(CTH_TYPES)}"
) )
@@ -112,12 +101,6 @@ def parse_cth_comment(body: str) -> dict[str, Any] | None:
fields[key] = match.group(2).strip() fields[key] = match.group(2).strip()
return { return {
"cth_type": cth_type, "cth_type": cth_type,
# The heading capture is unconstrained free text, so the parse states
# whether it satisfies the CTH_TYPES contract instead of leaving every
# reader to decide (or forget). Parsing stays total — an unknown type is
# still parsed and reported, never raised on — but a reader that turns
# the type into a durable value can now tell the two apart.
"cth_type_known": is_known_cth_type(cth_type),
"fields": fields, "fields": fields,
"raw_body": text, "raw_body": text,
} }
@@ -136,7 +119,7 @@ def assess_cth_comment(body: str) -> dict[str, Any]:
} }
cth_type = parsed.get("cth_type") or "" cth_type = parsed.get("cth_type") or ""
if not is_known_cth_type(cth_type): if cth_type not in CTH_TYPES:
reasons.append( reasons.append(
f"unknown CTH type '{cth_type}'; expected one of {sorted(CTH_TYPES)}" f"unknown CTH type '{cth_type}'; expected one of {sorted(CTH_TYPES)}"
) )
+13 -125
View File
@@ -65,6 +65,8 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| `/api/prompts` | JSON prompt export with workflow hashes | | `/api/prompts` | JSON prompt export with workflow hashes |
| `/runtime` | MCP runtime health and stale detection (#430) | | `/runtime` | MCP runtime health and stale detection (#430) |
| `/api/runtime` | JSON runtime health export | | `/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) | | `/audit` | Report audit paste + validator preview (#431) |
| `/api/audit` | JSON validator preview (POST `report_text`, optional `task_kind`) | | `/api/audit` | JSON validator preview (POST `report_text`, optional `task_kind`) |
| `/worktrees` | Worktree hygiene dashboard (#432) | | `/worktrees` | Worktree hygiene dashboard (#432) |
@@ -74,11 +76,6 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| `/api/actions/{id}/preview` | Mutation ledger preview (GET, read-only) | | `/api/actions/{id}/preview` | Mutation ledger preview (GET, read-only) |
| `/leases` | Lease and collision visibility (#433) | | `/leases` | Lease and collision visibility (#433) |
| `/api/leases` | JSON lease/collision export | | `/api/leases` | JSON lease/collision export |
| `/sessions` | Phase 1 shell stub — session inventory (backed by #636) |
| `/inventory` | Phase 1 shell stub — unified inventory (backed by #636) |
| `/timeline` | Phase 1 shell stub — workflow event timeline |
| `/policy` | Phase 1 shell stub — capability/role policy placeholder |
| `/insights` | Phase 1 shell stub — operational insights placeholder |
Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for `read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
@@ -238,25 +235,18 @@ health, workflow/schema SHA-256 hashes, and stale-runtime warnings when the
checkout is behind merged safety-gate changes. Restart guidance links to #420; checkout is behind merged safety-gate changes. Restart guidance links to #420;
no tokens or MCP restart actions are exposed. no tokens or MCP restart actions are exposed.
## Application shell — Phase 1 (#638) ## Policy & guardrail visibility (#646)
The console shell (`webui/layout.py`) renders a grouped navigation driven by a `/policy` (HTML) and `/api/v1/policy` (JSON) surface a **read-only** projection
single nav-config module, `webui/nav.py`. Nav groups follow the epic #631 of the major workflow guardrails — role separation/RBAC, lease lifecycle,
Phase 1 information architecture: **Health, Traffic, Runtime/Sessions, author worktree binding, merge confirmation, secret redaction, contamination
Projects, Inventory, Timeline, Policy** (placeholder), and **Insights** containment, allocator policy, audit logging, and mutation gating. Each entry
(placeholder). Live views and Phase 1 placeholders (`stub`) are declared in one carries source pointers to the file/module/doc that owns it, a compact active
place so the layout and the route table cannot drift. 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
The header carries two read-only status badges — an **environment** badge through the console redaction pass before it is emitted, so a planted or
(`local` for loopback binds, `remote` otherwise, derived from `WEBUI_HOST`) and accidental secret degrades to the placeholder rather than reaching a client.
a **mode: read-only** badge — plus a **Docs** link to this document. No The view never edits policy and exposes no gate-weakening toggle.
privileged action controls are present in the Phase 1 shell.
Not-yet-implemented surfaces (`/sessions`, `/inventory`, `/timeline`,
`/policy`, `/insights`) resolve to graceful read-only stub pages instead of
404s; their backing views land in later child issues of #631 (the inventory
surfaces are backed by #636). Mutating methods on stub routes still fail closed
with `read-only-mvp`.
## Deployment boundary (#435) ## Deployment boundary (#435)
@@ -317,108 +307,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; checkout is behind merged safety-gate changes. Restart guidance links to #420;
no tokens or MCP restart actions are exposed. 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.
### Filter authority, and refusing what cannot be answered
A filter dimension is only meaningful for a source whose records carry it.
Each source declares its own support in `_SOURCE_FILTER_SUPPORT` and reports it
per response as `supported_filters` / `unsupported_filters`:
| Source | issue | pr | session |
|---|---|---|---|
| `control_plane` | yes | yes | **no** — the `events` table is `(event_id, work_item_id, event_type, message, created_at)` and records no session |
| `gitea_handoff` | yes | yes | yes — a CTH comment declares its own `Session:` field |
`session_id` is read only from that declared CTH field. It is never inferred
from a work item, an actor, or message text, and a value that is
redaction-altering or bare-secret-shaped is dropped rather than emitted.
When **no source that ran** can carry a requested dimension, the request is
refused rather than answered: the response is `422` with `ok:false` and a
structured `error` naming `unsupported_filters` and the per-source reason. A
`200` with zero events would tell an operator that no such activity exists,
which is a stronger — and false — claim than "this cannot be answered here".
A source that *can* answer the dimension and simply matched nothing still
returns `200` with `ok:true` and an empty page.
### Redaction
Every free-text field (event messages, decision/proof text, roles, actors) 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.
Redaction also runs *before* any structured value is derived from free text.
`evidence_refs` are extracted from already-redacted proof/decision text, and a
commit reference is recognised only where the text declares one (`commit`,
`head`, `base`, `sha`, …). An undeclared 40-character hex run has the exact
shape of a Gitea access token, so it is never lifted out of prose into a
structured field. Every reference is then independently revalidated against an
allowed shape and a second redaction pass immediately before serialization;
anything unproven is dropped and the event is flagged `sensitive`.
### Tests
```bash
pytest tests/test_webui_timeline.py -q
```
## Tests ## Tests
```bash ```bash
+99 -135
View File
@@ -11242,163 +11242,127 @@ def gitea_reconcile_merged_cleanups(
if dry_run: if dry_run:
report["dry_run"] = True report["dry_run"] = True
report["executed"] = False report["executed"] = False
# #851: surface planned lifecycle order so dry-run matches execute.
report["planned_execution_orders"] = {
str(entry.get("pr_number")): entry.get("planned_execution_order") or []
for entry in (report.get("entries") or [])
}
return {"success": True, "performed": False, **report} return {"success": True, "performed": False, **report}
verify_preflight_purity( verify_preflight_purity(
remote, task="reconcile_merged_cleanups", org=org, repo=repo remote, task="reconcile_merged_cleanups", org=org, repo=repo
) )
actions: list[dict] = [] actions: list[dict] = []
project_root = _canonical_local_git_root()
def _ownership_records_for_branch(
head_branch: str, pr_num_int: int | None
) -> list[dict]:
ownership_bundle = _collect_branch_ownership_records(
remote=remote,
host=h,
org=o,
repo=r,
branch=head_branch,
pr_number=pr_num_int,
project_root=project_root,
auth=auth,
base_api=base,
)
ownership_records = list(ownership_bundle.get("records") or [])
if ownership_bundle.get("inventory_error"):
ownership_records.append(
{
"category": (
branch_cleanup_guard.OWNERSHIP_CATEGORY_INVENTORY_ERROR
),
"status": "unknown",
"remote": remote,
"host": h,
"org": o,
"repo": r,
"branch": head_branch,
"reclaim_allowed": False,
"role": "inventory",
}
)
return ownership_records
def _attempt_owned_remote_delete(
*,
head_branch: str,
pr_num_int: int | None,
after_worktree_removal: bool = False,
) -> dict:
"""Fail-closed remote delete with live ownership reassessment (#851)."""
import urllib.parse
ownership_records = _ownership_records_for_branch(head_branch, pr_num_int)
ownership = branch_cleanup_guard.assess_active_branch_ownership(
remote=remote,
org=o,
repo=r,
branch=head_branch,
host=h,
records=ownership_records,
)
if ownership.get("block"):
return {
"action": "delete_remote_branch",
"branch": head_branch,
"success": False,
"performed": False,
"delete_acknowledged": False,
"verified_absent": False,
"blocker_kind": "active_branch_ownership",
"reasons": ownership.get("reasons") or [],
"blocking_categories": ownership.get("blocking_categories") or [],
"after_worktree_removal": after_worktree_removal,
"ownership_reassessed": after_worktree_removal,
}
encoded = urllib.parse.quote(head_branch, safe="")
url = f"{base}/branches/{encoded}"
with _audited(
"delete_branch",
host=h,
remote=remote,
org=o,
repo=r,
target_branch=head_branch,
request_metadata={
"branch": head_branch,
"source": "reconcile_merged_cleanups",
"ownership_checked": True,
"after_worktree_removal": after_worktree_removal,
},
):
api_request("DELETE", url, auth)
readback = _probe_remote_branch(h, o, r, auth, head_branch)
readback_assessment = branch_cleanup_guard.assess_post_delete_readback(
readback
)
verified = bool(readback_assessment.get("verified_absent"))
return {
"action": "delete_remote_branch",
"branch": head_branch,
"success": bool(readback_assessment.get("ok")),
"performed": True,
"delete_acknowledged": True,
"verified_absent": verified,
"readback": readback_assessment.get("readback"),
"reasons": readback_assessment.get("reasons") or [],
"after_worktree_removal": after_worktree_removal,
"ownership_reassessed": after_worktree_removal,
}
for entry in report.get("entries") or []: for entry in report.get("entries") or []:
head_branch = entry.get("head_branch") or "" head_branch = entry.get("head_branch") or ""
remote_assessment = entry.get("remote_branch") or {} remote_assessment = entry.get("remote_branch") or {}
local_assessment = entry.get("local_worktree") or {} local_assessment = entry.get("local_worktree") or {}
pr_num = entry.get("pr_number")
try:
pr_num_int = int(pr_num) if pr_num is not None else None
except (TypeError, ValueError):
pr_num_int = None
# #851 lifecycle: when the target worktree is independently safe, remove if remote_assessment.get("safe_to_delete_remote"):
# it first so worktree_binding ownership does not permanently strand import urllib.parse
# both the worktree and the remote branch. Never skip worktree removal
# merely because remote delete would be blocked by that binding. pr_num = entry.get("pr_number")
# Ownership protection for remote delete remains fail-closed below. try:
worktree_removed = False pr_num_int = int(pr_num) if pr_num is not None else None
except (TypeError, ValueError):
pr_num_int = None
ownership_bundle = _collect_branch_ownership_records(
remote=remote,
host=h,
org=o,
repo=r,
branch=head_branch,
pr_number=pr_num_int,
project_root=_canonical_local_git_root(),
auth=auth,
base_api=base,
)
ownership_records = list(ownership_bundle.get("records") or [])
if ownership_bundle.get("inventory_error"):
ownership_records.append(
{
"category": (
branch_cleanup_guard.OWNERSHIP_CATEGORY_INVENTORY_ERROR
),
"status": "unknown",
"remote": remote,
"host": h,
"org": o,
"repo": r,
"branch": head_branch,
"reclaim_allowed": False,
"role": "inventory",
}
)
ownership = branch_cleanup_guard.assess_active_branch_ownership(
remote=remote,
org=o,
repo=r,
branch=head_branch,
host=h,
records=ownership_records,
)
if ownership.get("block"):
actions.append(
{
"action": "delete_remote_branch",
"branch": head_branch,
"success": False,
"performed": False,
"delete_acknowledged": False,
"verified_absent": False,
"blocker_kind": "active_branch_ownership",
"reasons": ownership.get("reasons") or [],
"blocking_categories": ownership.get(
"blocking_categories"
)
or [],
}
)
continue
encoded = urllib.parse.quote(head_branch, safe="")
url = f"{base}/branches/{encoded}"
with _audited(
"delete_branch",
host=h,
remote=remote,
org=o,
repo=r,
target_branch=head_branch,
request_metadata={
"branch": head_branch,
"source": "reconcile_merged_cleanups",
"ownership_checked": True,
},
):
api_request("DELETE", url, auth)
readback = _probe_remote_branch(h, o, r, auth, head_branch)
readback_assessment = branch_cleanup_guard.assess_post_delete_readback(
readback
)
verified = bool(readback_assessment.get("verified_absent"))
actions.append(
{
"action": "delete_remote_branch",
"branch": head_branch,
"success": bool(readback_assessment.get("ok")),
"performed": True,
"delete_acknowledged": True,
"verified_absent": verified,
"readback": readback_assessment.get("readback"),
"reasons": readback_assessment.get("reasons") or [],
}
)
if local_assessment.get("safe_to_remove_worktree"): if local_assessment.get("safe_to_remove_worktree"):
result = merged_cleanup_reconcile.remove_local_worktree( result = merged_cleanup_reconcile.remove_local_worktree(
project_root, _canonical_local_git_root(),
head_branch, head_branch,
worktree_path=local_assessment.get("worktree_path"), worktree_path=local_assessment.get("worktree_path"),
) )
actions.append({"action": "remove_local_worktree", **result}) actions.append({"action": "remove_local_worktree", **result})
# Idempotent resume: absent worktree is already gone.
msg = (result.get("message") or "").lower()
worktree_removed = bool(result.get("success")) or (
"not found" in msg
)
if remote_assessment.get("safe_to_delete_remote"):
actions.append(
_attempt_owned_remote_delete(
head_branch=head_branch,
pr_num_int=pr_num_int,
after_worktree_removal=worktree_removed,
)
)
for scratch in report.get("reviewer_scratch_entries") or []: for scratch in report.get("reviewer_scratch_entries") or []:
if not scratch.get("safe_to_remove_worktree"): if not scratch.get("safe_to_remove_worktree"):
continue continue
result = merged_cleanup_reconcile.remove_reviewer_scratch_worktree( result = merged_cleanup_reconcile.remove_reviewer_scratch_worktree(
project_root, scratch.get("worktree_path") or "" _canonical_local_git_root(), scratch.get("worktree_path") or ""
) )
actions.append({"action": "remove_reviewer_scratch_worktree", **result}) actions.append({"action": "remove_reviewer_scratch_worktree", **result})
-58
View File
@@ -566,10 +566,6 @@ def build_pr_cleanup_entry(
worktree_state=worktree_state, worktree_state=worktree_state,
active_lock=active_lock, active_lock=active_lock,
) )
planned = plan_cleanup_execution_order(
remote_assessment=remote,
local_assessment=local,
)
return { return {
"pr_number": pr_number, "pr_number": pr_number,
"issue_number": issue_number, "issue_number": issue_number,
@@ -580,63 +576,9 @@ def build_pr_cleanup_entry(
"merged": merged, "merged": merged,
"remote_branch": remote, "remote_branch": remote,
"local_worktree": local, "local_worktree": local,
# #851: dry-run and execute share the same lifecycle order description.
"planned_execution_order": planned,
} }
def plan_cleanup_execution_order(
*,
remote_assessment: dict[str, Any] | None,
local_assessment: dict[str, Any] | None,
) -> list[dict[str, Any]]:
"""Describe independent worktree-then-reassess-then-remote cleanup order (#851).
Remote ownership protection remains fail-closed at execute time. A worktree
that is independently safe to remove is never skipped merely because remote
deletion may be blocked by that same ``worktree_binding``.
"""
remote = remote_assessment or {}
local = local_assessment or {}
steps: list[dict[str, Any]] = []
worktree_safe = bool(local.get("safe_to_remove_worktree"))
remote_safe = bool(remote.get("safe_to_delete_remote"))
if worktree_safe:
steps.append(
{
"action": "remove_local_worktree",
"reason": "independently_safe_to_remove",
"phase": 1,
}
)
if remote_safe:
if worktree_safe:
steps.append(
{
"action": "reassess_branch_ownership",
"reason": "after_worktree_removal_clear_worktree_binding",
"phase": 2,
}
)
steps.append(
{
"action": "delete_remote_branch",
"reason": "only_if_independently_safe_after_reassessment",
"phase": 3,
}
)
else:
steps.append(
{
"action": "delete_remote_branch",
"reason": "safe_to_delete_and_no_independent_worktree_removal",
"phase": 1,
}
)
return steps
def build_reconciliation_report( def build_reconciliation_report(
*, *,
project_root: str, project_root: str,
-372
View File
@@ -1266,378 +1266,6 @@ class TestSecondRemediationIntegration(unittest.TestCase):
self.assertIn("delete_acknowledged", delete_actions[0]) self.assertIn("delete_acknowledged", delete_actions[0])
self.assertTrue(delete_actions[0].get("verified_absent")) self.assertTrue(delete_actions[0].get("verified_absent"))
def test_issue_851_worktree_removed_when_remote_blocked_only_by_worktree_binding(self):
"""#851: remote blocked by worktree_binding must not skip safe worktree removal.
Lifecycle: remove clean owned worktree → reassess ownership → delete
remote only if independently safe. Unrelated entries stay untouched.
"""
from mcp_server import gitea_reconcile_merged_cleanups
target_branch = "fix/issue-844-exclude-epic-containers"
foreign_branch = "fix/issue-999-unrelated-active"
worktree_path = "/tmp/branches/fix-issue-844-exclude-epic-containers"
ownership_calls = []
remove_calls = []
delete_api_calls = []
def fake_collect(**kwargs):
ownership_calls.append(dict(kwargs))
# Ownership is reassessed *after* independent worktree removal (#851).
# Target worktree is already gone → no worktree_binding remains.
# Foreign branch keeps an active author lease → remote delete blocked.
if kwargs.get("branch") == foreign_branch:
# Match session-bound org/repo + host used by the tool resolve path.
return {
"records": [
{
"category": guard.OWNERSHIP_CATEGORY_AUTHOR_LEASE,
"status": "active",
"remote": kwargs.get("remote") or "prgs",
"host": kwargs.get("host") or "gitea.example.com",
"org": kwargs.get("org") or "Scaled-Tech-Consulting",
"repo": kwargs.get("repo") or "Gitea-Tools",
"branch": foreign_branch,
"reclaim_allowed": False,
}
],
"inventory_error": False,
}
return {"records": [], "inventory_error": False}
def fake_remove(project_root, branch, worktree_path=None):
remove_calls.append(
{"branch": branch, "worktree_path": worktree_path}
)
return {
"success": True,
"performed": True,
"message": f"removed worktree {worktree_path}",
"worktree_path": worktree_path,
}
def fake_probe(h, o, r, auth, br):
return guard.classify_branch_readback_http_status(
404, not_found_scope=guard.NOT_FOUND_SCOPE_BRANCH
)
def fake_api(method, url, auth, **kwargs):
if method == "DELETE":
delete_api_calls.append(url)
return {}
report = {
"entries": [
{
"pr_number": 848,
"head_branch": target_branch,
"remote_branch": {"safe_to_delete_remote": True},
"local_worktree": {
"safe_to_remove_worktree": True,
"worktree_path": worktree_path,
},
},
{
"pr_number": 999,
"head_branch": foreign_branch,
"remote_branch": {"safe_to_delete_remote": True},
"local_worktree": {
"safe_to_remove_worktree": False,
"worktree_path": None,
},
},
],
"reviewer_scratch_entries": [],
}
patch(
"mcp_server.get_profile",
return_value={
"profile_name": "prgs-reconciler",
"role": "reconciler",
"allowed_operations": [
"gitea.read",
"gitea.branch.delete",
"gitea.pr.close",
],
"forbidden_operations": [],
},
).start()
patch("mcp_server.api_get_all", return_value=[]).start()
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
return_value=report,
).start()
patch(
"mcp_server.merged_cleanup_reconcile.discover_reviewer_scratch_worktrees",
return_value=[],
).start()
patch(
"mcp_server.audit_reconciliation_mode.check_cleanup_execution_allowed",
return_value=(True, []),
).start()
patch("mcp_server.verify_preflight_purity", return_value=None).start()
patch(
"mcp_server._collect_branch_ownership_records",
side_effect=fake_collect,
).start()
patch("mcp_server._probe_remote_branch", side_effect=fake_probe).start()
patch(
"mcp_server.merged_cleanup_reconcile.remove_local_worktree",
side_effect=fake_remove,
).start()
self.mock_api.side_effect = fake_api
res = gitea_reconcile_merged_cleanups(
dry_run=False,
execute_confirmed=True,
remote="prgs",
)
self.assertTrue(res.get("performed") or res.get("executed"))
actions = res.get("actions") or []
remove_actions = [
a for a in actions if a.get("action") == "remove_local_worktree"
]
self.assertEqual(len(remove_actions), 1, actions)
self.assertTrue(remove_actions[0].get("success"))
self.assertEqual(remove_calls[0]["branch"], target_branch)
self.assertEqual(remove_calls[0]["worktree_path"], worktree_path)
# Target remote delete succeeds after worktree removal + reassessment.
target_deletes = [
a
for a in actions
if a.get("action") == "delete_remote_branch"
and a.get("branch") == target_branch
]
self.assertEqual(len(target_deletes), 1, actions)
self.assertTrue(target_deletes[0].get("success"))
self.assertTrue(target_deletes[0].get("after_worktree_removal"))
self.assertTrue(target_deletes[0].get("ownership_reassessed"))
self.assertTrue(target_deletes[0].get("verified_absent"))
# Foreign branch remains protected (author lease) and is not deleted.
foreign_deletes = [
a
for a in actions
if a.get("action") == "delete_remote_branch"
and a.get("branch") == foreign_branch
]
self.assertEqual(len(foreign_deletes), 1, actions)
self.assertFalse(foreign_deletes[0].get("success"))
self.assertEqual(
foreign_deletes[0].get("blocker_kind"), "active_branch_ownership"
)
self.assertIn(
guard.OWNERSHIP_CATEGORY_AUTHOR_LEASE,
foreign_deletes[0].get("blocking_categories") or [],
)
# Only the target branch should hit the DELETE API.
self.assertEqual(len(delete_api_calls), 1)
# Ownership collected for target (post-removal) and foreign; worktree
# removal happened before target remote delete in the action log.
target_idx = next(
i
for i, a in enumerate(actions)
if a.get("action") == "remove_local_worktree"
)
delete_idx = next(
i
for i, a in enumerate(actions)
if a.get("action") == "delete_remote_branch"
and a.get("branch") == target_branch
and a.get("success")
)
self.assertLess(target_idx, delete_idx)
def test_issue_851_dirty_worktree_not_removed_and_remote_stays_protected(self):
"""#851: dirty/foreign worktrees remain protected; no unsafe cleanup."""
from mcp_server import gitea_reconcile_merged_cleanups
branch = "fix/issue-851-dirty"
remove_calls = []
def fake_collect(**kwargs):
return {
"records": [
{
"category": guard.OWNERSHIP_CATEGORY_WORKTREE_BINDING,
"status": "active",
"remote": kwargs.get("remote") or "prgs",
"host": kwargs.get("host") or "gitea.example.com",
"org": kwargs.get("org") or "Scaled-Tech-Consulting",
"repo": kwargs.get("repo") or "Gitea-Tools",
"branch": branch,
"reclaim_allowed": False,
}
],
"inventory_error": False,
}
report = {
"entries": [
{
"pr_number": 851,
"head_branch": branch,
"remote_branch": {"safe_to_delete_remote": True},
"local_worktree": {
"safe_to_remove_worktree": False,
"worktree_path": "/tmp/dirty-wt",
},
}
],
"reviewer_scratch_entries": [],
}
patch(
"mcp_server.get_profile",
return_value={
"profile_name": "prgs-reconciler",
"role": "reconciler",
"allowed_operations": [
"gitea.read",
"gitea.branch.delete",
],
"forbidden_operations": [],
},
).start()
patch("mcp_server.api_get_all", return_value=[]).start()
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
return_value=report,
).start()
patch(
"mcp_server.merged_cleanup_reconcile.discover_reviewer_scratch_worktrees",
return_value=[],
).start()
patch(
"mcp_server.audit_reconciliation_mode.check_cleanup_execution_allowed",
return_value=(True, []),
).start()
patch("mcp_server.verify_preflight_purity", return_value=None).start()
patch(
"mcp_server._collect_branch_ownership_records",
side_effect=fake_collect,
).start()
patch(
"mcp_server.merged_cleanup_reconcile.remove_local_worktree",
side_effect=lambda *a, **k: remove_calls.append(k) or {
"success": True,
"performed": True,
},
).start()
self.mock_api.side_effect = lambda *a, **k: {}
res = gitea_reconcile_merged_cleanups(
dry_run=False,
execute_confirmed=True,
remote="prgs",
)
actions = res.get("actions") or []
self.assertEqual(remove_calls, [])
self.assertFalse(
any(a.get("action") == "remove_local_worktree" for a in actions)
)
deletes = [
a for a in actions if a.get("action") == "delete_remote_branch"
]
self.assertEqual(len(deletes), 1)
self.assertFalse(deletes[0].get("success"))
self.assertEqual(deletes[0].get("blocker_kind"), "active_branch_ownership")
self.assertIn(
guard.OWNERSHIP_CATEGORY_WORKTREE_BINDING,
deletes[0].get("blocking_categories") or [],
)
def test_issue_851_idempotent_resume_when_worktree_already_absent(self):
"""#851: partial failures remain resumable and idempotent."""
from mcp_server import gitea_reconcile_merged_cleanups
branch = "fix/issue-851-resume"
ownership_calls = []
def fake_collect(**kwargs):
ownership_calls.append(kwargs)
return {"records": [], "inventory_error": False}
def fake_remove(project_root, branch, worktree_path=None):
return {
"success": False,
"performed": False,
"message": f"worktree not found: {worktree_path}",
}
def fake_probe(h, o, r, auth, br):
return guard.classify_branch_readback_http_status(
404, not_found_scope=guard.NOT_FOUND_SCOPE_BRANCH
)
report = {
"entries": [
{
"pr_number": 851,
"head_branch": branch,
"remote_branch": {"safe_to_delete_remote": True},
"local_worktree": {
"safe_to_remove_worktree": True,
"worktree_path": "/tmp/already-gone",
},
}
],
"reviewer_scratch_entries": [],
}
patch(
"mcp_server.get_profile",
return_value={
"profile_name": "prgs-reconciler",
"role": "reconciler",
"allowed_operations": [
"gitea.read",
"gitea.branch.delete",
],
"forbidden_operations": [],
},
).start()
patch("mcp_server.api_get_all", return_value=[]).start()
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
return_value=report,
).start()
patch(
"mcp_server.merged_cleanup_reconcile.discover_reviewer_scratch_worktrees",
return_value=[],
).start()
patch(
"mcp_server.audit_reconciliation_mode.check_cleanup_execution_allowed",
return_value=(True, []),
).start()
patch("mcp_server.verify_preflight_purity", return_value=None).start()
patch(
"mcp_server._collect_branch_ownership_records",
side_effect=fake_collect,
).start()
patch("mcp_server._probe_remote_branch", side_effect=fake_probe).start()
patch(
"mcp_server.merged_cleanup_reconcile.remove_local_worktree",
side_effect=fake_remove,
).start()
self.mock_api.side_effect = lambda *a, **k: {}
res = gitea_reconcile_merged_cleanups(
dry_run=False,
execute_confirmed=True,
remote="prgs",
)
actions = res.get("actions") or []
removes = [a for a in actions if a.get("action") == "remove_local_worktree"]
deletes = [a for a in actions if a.get("action") == "delete_remote_branch"]
self.assertEqual(len(removes), 1)
self.assertFalse(removes[0].get("success"))
self.assertEqual(len(deletes), 1)
self.assertTrue(deletes[0].get("success"))
self.assertTrue(deletes[0].get("after_worktree_removal"))
self.assertTrue(ownership_calls)
if __name__ == "__main__": if __name__ == "__main__":
-53
View File
@@ -12,59 +12,6 @@ import merged_cleanup_reconcile as mcr # noqa: E402
class TestMergedCleanupAssessment(unittest.TestCase): class TestMergedCleanupAssessment(unittest.TestCase):
def test_issue_851_plan_order_worktree_then_reassess_then_remote(self):
"""#851 dry-run plan: remove worktree, reassess ownership, then remote."""
plan = mcr.plan_cleanup_execution_order(
remote_assessment={"safe_to_delete_remote": True},
local_assessment={"safe_to_remove_worktree": True},
)
actions = [s["action"] for s in plan]
self.assertEqual(
actions,
[
"remove_local_worktree",
"reassess_branch_ownership",
"delete_remote_branch",
],
)
self.assertEqual(plan[0]["phase"], 1)
self.assertEqual(plan[-1]["phase"], 3)
self.assertIn("independently_safe", plan[0]["reason"])
self.assertIn("reassessment", plan[-1]["reason"])
def test_issue_851_plan_remote_only_when_worktree_not_safe(self):
plan = mcr.plan_cleanup_execution_order(
remote_assessment={"safe_to_delete_remote": True},
local_assessment={"safe_to_remove_worktree": False},
)
self.assertEqual([s["action"] for s in plan], ["delete_remote_branch"])
self.assertNotIn("reassess_branch_ownership", [s["action"] for s in plan])
def test_issue_851_plan_worktree_only_when_remote_not_safe(self):
plan = mcr.plan_cleanup_execution_order(
remote_assessment={"safe_to_delete_remote": False},
local_assessment={"safe_to_remove_worktree": True},
)
self.assertEqual([s["action"] for s in plan], ["remove_local_worktree"])
def test_issue_851_entry_includes_planned_execution_order(self):
entry = mcr.build_pr_cleanup_entry(
pr={
"number": 848,
"title": "Closes #844",
"body": "",
"merged_at": "2026-07-23T00:00:00Z",
"head": {"ref": "fix/issue-844-x", "sha": "a" * 40},
},
project_root="/tmp/not-a-real-root",
open_pr_heads=set(),
remote_branch_exists=True,
head_on_master=True,
delete_capability_allowed=True,
)
self.assertIn("planned_execution_order", entry)
self.assertIsInstance(entry["planned_execution_order"], list)
def test_extract_linked_issue_from_closes(self): def test_extract_linked_issue_from_closes(self):
issue = mcr.extract_linked_issue( issue = mcr.extract_linked_issue(
"feat: cleanup (Closes #269)", "feat: cleanup (Closes #269)",
+222
View File
@@ -0,0 +1,222 @@
"""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()
-135
View File
@@ -1,135 +0,0 @@
"""Tests for the Phase 1 operator console application shell (#638)."""
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.routing import Route
from starlette.testclient import TestClient
from webui import layout
from webui.app import create_app
from webui.nav import NAV_GROUPS, STUB_PAGES, nav_hrefs
class TestShellNav(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_nav_group_labels_present(self):
text = self.client.get("/").text
for group in NAV_GROUPS:
with self.subTest(group=group.label):
self.assertIn(f">{group.label}<", text)
def test_phase1_group_labels_cover_expected_ia(self):
labels = {group.label for group in NAV_GROUPS}
for expected in (
"Health",
"Traffic",
"Runtime/Sessions",
"Projects",
"Inventory",
"Timeline",
"Policy",
"Insights",
):
with self.subTest(label=expected):
self.assertIn(expected, labels)
def test_every_nav_href_resolves_to_a_get_route(self):
app = create_app()
get_paths = {
route.path
for route in app.routes
if isinstance(route, Route) and "GET" in route.methods
}
for href in nav_hrefs():
with self.subTest(href=href):
self.assertIn(href, get_paths, f"nav href {href} has no GET route")
def test_legacy_hrefs_still_navigable(self):
text = self.client.get("/").text
for href in ("/queue", "/projects", "/prompts", "/runtime",
"/audit", "/worktrees", "/leases", "/actions"):
with self.subTest(href=href):
self.assertIn(f'href="{href}"', text)
class TestShellBadges(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_mode_badge_present(self):
self.assertIn("mode: read-only", self.client.get("/").text)
def test_environment_badge_present(self):
self.assertIn("env:", self.client.get("/").text)
def test_default_environment_is_local(self):
self.assertEqual(layout.environment_label(), "local")
def test_remote_bind_reports_remote_environment(self):
import os
prior = os.environ.get("WEBUI_HOST")
os.environ["WEBUI_HOST"] = "10.0.0.5"
try:
self.assertEqual(layout.environment_label(), "remote")
finally:
if prior is None:
os.environ.pop("WEBUI_HOST", None)
else:
os.environ["WEBUI_HOST"] = prior
def test_docs_link_present(self):
text = self.client.get("/").text
self.assertIn(layout.DOCS_URL, text)
self.assertIn(">Docs<", text)
class TestShellStubs(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_stub_routes_render_200(self):
for path, (title, _desc) in STUB_PAGES.items():
with self.subTest(path=path):
response = self.client.get(path)
self.assertEqual(response.status_code, 200, path)
self.assertIn(title, response.text)
self.assertIn("placeholder", response.text)
def test_stub_routes_are_read_only(self):
for path in STUB_PAGES:
with self.subTest(path=path):
response = self.client.post(path)
self.assertEqual(response.status_code, 405)
self.assertEqual(response.json()["error"], "read-only-mvp")
def test_stub_pages_carry_nav_and_badges(self):
response = self.client.get("/inventory")
self.assertIn("mode: read-only", response.text)
self.assertIn('href="/queue"', response.text)
class TestShellHome(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_home_summarizes_console(self):
text = self.client.get("/").text
self.assertIn("Operator console", text)
self.assertIn("Phase 1", text)
def test_home_links_legacy_pages(self):
text = self.client.get("/").text
self.assertIn("MVP legacy pages", text)
for href in ("/queue", "/audit", "/leases"):
with self.subTest(href=href):
self.assertIn(f'href="{href}"', text)
if __name__ == "__main__":
unittest.main()
File diff suppressed because it is too large Load Diff
+27 -154
View File
@@ -12,7 +12,6 @@ from starlette.routing import Route
from webui.deployment_boundary import deployment_snapshot from webui.deployment_boundary import deployment_snapshot
from webui.layout import render_page from webui.layout import render_page
from webui.nav import NAV_GROUPS, STUB_PAGES
from webui.project_registry import ( from webui.project_registry import (
ProjectRegistry, ProjectRegistry,
RegistryError, RegistryError,
@@ -46,7 +45,8 @@ from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as wo
from webui.worktree_views import render_worktrees_page from webui.worktree_views import render_worktrees_page
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
from webui.runtime_views import render_runtime_page from webui.runtime_views import render_runtime_page
from webui.timeline import load_timeline, snapshot_to_dict as timeline_snapshot_to_dict 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.system_health import ( from webui.system_health import (
API_PATH as SYSTEM_HEALTH_API_PATH, API_PATH as SYSTEM_HEALTH_API_PATH,
load_system_health, load_system_health,
@@ -67,62 +67,25 @@ def _stub_page(title: str, description: str) -> HTMLResponse:
return HTMLResponse(render_page(title=title, body_html=body)) return HTMLResponse(render_page(title=title, body_html=body))
_LEGACY_PAGES = (
("/queue", "Queue", "live PR and issue dashboard (#429)"),
("/projects", "Projects", "registry and onboarding (#427)"),
("/prompts", "Prompts", "canonical workflow prompt library (#428)"),
("/runtime", "Runtime", "MCP health and stale-runtime detection (#430)"),
("/audit", "Audit", "final-report paste and validator preview (#431)"),
("/worktrees", "Worktrees", "branch hygiene dashboard (#432)"),
("/leases", "Leases", "collision and lease visibility (#433)"),
("/actions", "Actions", "gated write-action framework (#434)"),
)
def _render_home_nav_groups() -> str:
groups = []
for group in NAV_GROUPS:
items = "".join(
f'<li><a href="{item.href}">{item.label}</a>'
+ ("" if item.status == "live" else " <span class=\"muted\">(stub)</span>")
+ "</li>"
for item in group.items
)
groups.append(f"<h3>{group.label}</h3><ul>{items}</ul>")
return "".join(groups)
async def home(_request: Request) -> HTMLResponse: async def home(_request: Request) -> HTMLResponse:
legacy = "".join(
f"<li><strong>{label}</strong> — {desc} "
f'(<a href="{href}">{href}</a>)</li>'
for href, label, desc in _LEGACY_PAGES
)
body = ( body = (
"<h2>Operator console</h2>" "<h2>Operator console</h2>"
"<p>Read-only home for the MCP Control Plane Phase 1 operator console. " "<p>Local entry point for MCP Control Plane operational views.</p>"
"Gitea, MCP capability gates, and canonical workflows remain the source " "<ul>"
"of truth; this console never mutates them.</p>" "<li><strong>Queue</strong> — live PR and issue dashboard (#429)</li>"
"<h2>Phase 1 surfaces</h2>" "<li><strong>Projects</strong> — registry and onboarding (#427)</li>"
+ _render_home_nav_groups() "<li><strong>Prompts</strong> — canonical workflow prompt library (#428)</li>"
+ "<h2>MVP legacy pages</h2>" "<li><strong>Runtime</strong> — MCP health and stale-runtime detection (#430)</li>"
"<ul>" + legacy + "</ul>" "<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>"
"<li><strong>Actions</strong> — gated write-action framework (#434)</li>"
"</ul>"
) )
return HTMLResponse(render_page(title="Home", body_html=body)) return HTMLResponse(render_page(title="Home", body_html=body))
async def phase_stub(request: Request) -> HTMLResponse:
"""Graceful read-only placeholder for a not-yet-implemented Phase 1 surface."""
title, description = STUB_PAGES[request.url.path]
body = (
f"<h2>{title}</h2>"
f'<div class="stub"><p>{description}</p>'
"<p>Phase 1 shell placeholder — no write actions. Tracked under "
"epic #631.</p></div>"
)
return HTMLResponse(render_page(title=title, body_html=body))
async def health(_request: Request) -> JSONResponse: async def health(_request: Request) -> JSONResponse:
"""Liveness only — deliberately cheap, runs no dependency probe (#634). """Liveness only — deliberately cheap, runs no dependency probe (#634).
@@ -283,6 +246,17 @@ async def api_runtime(_request: Request) -> JSONResponse:
return JSONResponse(runtime_snapshot_to_dict(load_runtime_snapshot())) 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]: async def _parse_audit_form(request: Request) -> tuple[str, str | None]:
if request.method == "GET": if request.method == "GET":
return "", None return "", None
@@ -450,104 +424,6 @@ 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,
)
# A filter no surviving source can carry is refused, not answered empty:
# a 200 with zero events would tell the operator no such activity exists.
status_code = 200 if snapshot.ok else 422
return JSONResponse(timeline_snapshot_to_dict(snapshot), status_code=status_code)
async def method_not_allowed(request: Request, _exc: Exception) -> Response: async def method_not_allowed(request: Request, _exc: Exception) -> Response:
path = request.url.path path = request.url.path
if path in _AUDIT_MUTATION_PATHS and request.method == "POST": if path in _AUDIT_MUTATION_PATHS and request.method == "POST":
@@ -587,7 +463,8 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/api/prompts", api_prompts, methods=["GET"]), Route("/api/prompts", api_prompts, methods=["GET"]),
Route("/runtime", runtime, methods=["GET"]), Route("/runtime", runtime, methods=["GET"]),
Route("/api/runtime", api_runtime, methods=["GET"]), Route("/api/runtime", api_runtime, methods=["GET"]),
Route("/api/v1/timeline", api_v1_timeline, methods=["GET"]), Route("/policy", policy, methods=["GET"]),
Route("/api/v1/policy", api_v1_policy, methods=["GET"]),
Route("/audit", audit, methods=["GET", "POST"]), Route("/audit", audit, methods=["GET", "POST"]),
Route("/api/audit", api_audit, methods=["GET", "POST"]), Route("/api/audit", api_audit, methods=["GET", "POST"]),
Route("/worktrees", worktrees, methods=["GET"]), Route("/worktrees", worktrees, methods=["GET"]),
@@ -611,10 +488,6 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
api_console_security_model, api_console_security_model,
methods=["GET"], methods=["GET"],
), ),
*[
Route(path, phase_stub, methods=["GET"])
for path in STUB_PAGES
],
], ],
exception_handlers={405: method_not_allowed}, exception_handlers={405: method_not_allowed},
) )
+18 -95
View File
@@ -2,66 +2,29 @@
from __future__ import annotations from __future__ import annotations
import os NAV_ITEMS = (
("/", "Home"),
from webui.nav import NAV_GROUPS ("/queue", "Queue"),
("/projects", "Projects"),
("/prompts", "Prompts"),
("/runtime", "Runtime"),
("/policy", "Policy"),
("/audit", "Audit"),
("/worktrees", "Worktrees"),
("/leases", "Leases"),
("/actions", "Actions"),
)
MVP_NOTICE = ( MVP_NOTICE = (
"Read-only MVP — Gitea, MCP tools, and canonical workflows remain the " "Read-only MVP — Gitea, MCP tools, and canonical workflows remain the "
"source of truth. No mutation endpoints." "source of truth. No mutation endpoints."
) )
# Canonical docs entry point surfaced from the shell header (#638).
DOCS_URL = (
"https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/src/branch/"
"master/docs/webui-local-dev.md"
)
_LOCAL_HOSTS = frozenset({"", "127.0.0.1", "localhost", "::1"})
def environment_label() -> str:
"""Classify the serving environment as ``local`` or ``remote`` (#638).
Derived from the same ``WEBUI_HOST`` default the app binds to; loopback
hosts are ``local``, anything else is ``remote``. Read-only signal only.
"""
host = (os.environ.get("WEBUI_HOST", "127.0.0.1") or "").strip().lower()
return "local" if host in _LOCAL_HOSTS else "remote"
def _render_nav() -> str:
groups_html = []
for group in NAV_GROUPS:
links = "".join(
f'<a href="{item.href}"'
+ (' class="nav-stub"' if item.status == "stub" else "")
+ f'>{item.label}</a>'
for item in group.items
)
groups_html.append(
'<div class="nav-group">'
f'<span class="nav-group-label">{group.label}</span>'
f'<span class="nav-group-links">{links}</span>'
"</div>"
)
return "".join(groups_html)
def _render_badges() -> str:
env = environment_label()
return (
'<div class="header-badges">'
f'<span class="badge env-badge env-{env}">env: {env}</span>'
'<span class="badge mode-badge">mode: read-only</span>'
f'<a class="badge docs-link" href="{DOCS_URL}">Docs</a>'
"</div>"
)
def render_page(*, title: str, body_html: str, extra_head: str = "") -> str: def render_page(*, title: str, body_html: str, extra_head: str = "") -> str:
nav_links = _render_nav() nav_links = "".join(
header_badges = _render_badges() f'<a href="{href}">{label}</a>' for href, label in NAV_ITEMS
)
return f"""<!DOCTYPE html> return f"""<!DOCTYPE html>
<html lang="en"> <html lang="en">
<head> <head>
@@ -91,58 +54,21 @@ def render_page(*, title: str, body_html: str, extra_head: str = "") -> str:
padding: 0.75rem 1.25rem; padding: 0.75rem 1.25rem;
}} }}
header h1 {{ header h1 {{
margin: 0; margin: 0 0 0.5rem;
font-size: 1.1rem; font-size: 1.1rem;
font-weight: 600; font-weight: 600;
}} }}
.header-top {{
display: flex;
flex-wrap: wrap;
align-items: center;
justify-content: space-between;
gap: 0.5rem 1rem;
margin-bottom: 0.6rem;
}}
.header-badges {{ display: inline-flex; flex-wrap: wrap; gap: 0.4rem; }}
.env-badge.env-local {{ color: #8fd19e; border-color: #3d6b4a; }}
.env-badge.env-remote {{ color: #e0c27a; border-color: #6b5730; }}
.mode-badge {{ color: #9ec8f0; border-color: #3d5f7a; }}
a.docs-link {{
color: var(--accent);
border-color: var(--accent);
text-decoration: none;
text-transform: none;
}}
a.docs-link:hover {{ filter: brightness(1.12); }}
nav {{ nav {{
display: flex; display: flex;
flex-wrap: wrap; flex-wrap: wrap;
gap: 0.5rem 1.25rem; gap: 0.75rem 1rem;
}} }}
.nav-group {{
display: flex;
flex-direction: column;
gap: 0.15rem;
}}
.nav-group-label {{
font-size: 0.68rem;
text-transform: uppercase;
letter-spacing: 0.04em;
color: var(--muted);
}}
.nav-group-links {{ display: inline-flex; flex-wrap: wrap; gap: 0.6rem; }}
nav a {{ nav a {{
color: var(--accent); color: var(--accent);
text-decoration: none; text-decoration: none;
font-size: 0.9rem; font-size: 0.9rem;
}} }}
nav a:hover {{ text-decoration: underline; }} nav a:hover {{ text-decoration: underline; }}
nav a.nav-stub {{ color: var(--muted); }}
nav a.nav-stub::after {{
content: " ·stub";
font-size: 0.7rem;
color: var(--muted);
}}
main {{ main {{
max-width: 52rem; max-width: 52rem;
margin: 0 auto; margin: 0 auto;
@@ -241,10 +167,7 @@ def render_page(*, title: str, body_html: str, extra_head: str = "") -> str:
</head> </head>
<body> <body>
<header> <header>
<div class="header-top"> <h1>MCP Control Plane</h1>
<h1>MCP Control Plane</h1>
{header_badges}
</div>
<nav>{nav_links}</nav> <nav>{nav_links}</nav>
</header> </header>
<main> <main>
-111
View File
@@ -1,111 +0,0 @@
"""Navigation IA for the Phase 1 operator console shell (#638).
Single source of truth for the console navigation so ``webui/layout.py`` and
the ``webui/app.py`` route table stay aligned with epic #631. Read-only: every
destination is a GET view or a Phase 1 placeholder. No mutation links.
Nav groups follow the #631 Phase 1 information architecture: Health, Traffic,
Runtime/Sessions, Projects, Inventory, Timeline, Policy (placeholder), and
Insights (placeholder). Later-phase surfaces are declared as ``stub`` items and
backed by ``STUB_PAGES`` so their nav links resolve to a graceful placeholder
instead of a 404.
"""
from __future__ import annotations
from dataclasses import dataclass
@dataclass(frozen=True)
class NavItem:
"""A single navigation destination.
``status`` is ``"live"`` for implemented views and ``"stub"`` for Phase 1
placeholders whose backing view lands in a later child issue.
"""
href: str
label: str
status: str = "live"
@dataclass(frozen=True)
class NavGroup:
label: str
items: tuple[NavItem, ...]
NAV_GROUPS: tuple[NavGroup, ...] = (
NavGroup("Health", (
NavItem("/health", "Liveness"),
)),
NavGroup("Traffic", (
NavItem("/queue", "Queue"),
NavItem("/leases", "Leases"),
NavItem("/actions", "Actions"),
)),
NavGroup("Runtime/Sessions", (
NavItem("/runtime", "Runtime health"),
NavItem("/sessions", "Sessions", "stub"),
)),
NavGroup("Projects", (
NavItem("/projects", "Projects"),
)),
NavGroup("Inventory", (
NavItem("/inventory", "Inventory", "stub"),
NavItem("/worktrees", "Worktrees"),
)),
NavGroup("Timeline", (
NavItem("/timeline", "Timeline", "stub"),
)),
NavGroup("Policy", (
NavItem("/policy", "Policy", "stub"),
NavItem("/prompts", "Prompts"),
)),
NavGroup("Insights", (
NavItem("/insights", "Insights", "stub"),
NavItem("/audit", "Audit"),
)),
)
# Phase 1 placeholder destinations whose backing views land in later child
# issues of epic #631. Each maps a path to (title, description). Routes are
# registered so nav links resolve to a graceful, read-only stub page.
STUB_PAGES: dict[str, tuple[str, str]] = {
"/sessions": (
"Sessions",
"Active session, capability, and role inventory. Backed by the unified "
"inventory API (#636) once it lands.",
),
"/inventory": (
"Inventory",
"Unified sessions, leases, locks, namespaces, and worktree inventory. "
"Backed by the Phase 1 inventory API (#636).",
),
"/timeline": (
"Timeline",
"Workflow event timeline across issues and PRs. A later Phase 1 surface.",
),
"/policy": (
"Policy",
"Capability and role policy surface. Placeholder until a later phase.",
),
"/insights": (
"Insights",
"Aggregate operational insights and trends. Placeholder until a later "
"phase.",
),
}
def iter_nav_items():
"""Yield every ``NavItem`` across all groups in declared order."""
for group in NAV_GROUPS:
for item in group.items:
yield item
def nav_hrefs() -> tuple[str, ...]:
"""Return every navigation href in declared order."""
return tuple(item.href for item in iter_nav_items())
+387
View File
@@ -0,0 +1,387 @@
"""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
@@ -0,0 +1,104 @@
"""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>"
)
-906
View File
@@ -1,906 +0,0 @@
"""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.
- **Answerable filters only.** Each source declares which filter dimensions it
can actually answer. A filter dimension no source that ran can carry is
refused with an explicit reason rather than silently matching nothing: an
empty page from an unanswerable filter reads to an operator as "no such
activity", which is a different — and false — statement.
- **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, and *before* any
structured value is derived from it — evidence references are extracted from
redacted text, then independently revalidated before serialization. An
unredactable value becomes the placeholder, and a value that cannot be proven
safe is dropped — 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, replace
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,
}
# The filter dimensions the query layer accepts.
FILTER_ISSUE = "issue"
FILTER_PR = "pr"
FILTER_SESSION = "session"
# Which dimensions each source can actually answer. This is a property of the
# underlying records, not of the query code: the control-plane ``events`` table
# is (event_id, work_item_id, event_type, message, created_at) and carries no
# session identity at all, so no control-plane event can ever match a session
# filter. A CTH handoff comment can declare its session as a field, so the
# handoff source answers all three. Filtering on a dimension the surviving
# sources cannot carry is refused in ``load_timeline`` rather than answered
# with an empty page.
_SOURCE_FILTER_SUPPORT: dict[str, tuple[str, ...]] = {
SOURCE_CONTROL_PLANE: (FILTER_ISSUE, FILTER_PR),
SOURCE_GITEA_HANDOFF: (FILTER_ISSUE, FILTER_PR, FILTER_SESSION),
}
# Why a source cannot answer a dimension, for the refusal reason an operator reads.
_SOURCE_FILTER_LIMITS: dict[tuple[str, str], str] = {
(SOURCE_CONTROL_PLANE, FILTER_SESSION): (
"control-plane events carry no session identity "
"(the events table has no session column)"
),
}
# 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)
# A commit reference is only recognised when the text *declares* it as one.
# A bare lowercase hex run is not evidence of anything: at 40 characters it is
# exactly the shape of a Gitea personal access token, and at 7 it also matches
# ordinary words such as "defaced". Requiring an anchoring keyword keeps real
# references ("commit abc1234", "at head a209756...", "base caaae9b6") usable
# while refusing to lift an undeclared secret-shaped run out of free text.
_SHA_RE = re.compile(
r"(?i:\b(?:commit|sha|head|base|parent|revision|rev|merge[- ]base)\b[\s:=@#]*)"
r"([0-9a-f]{7,40})\b"
)
# Shapes a serialized evidence reference is allowed to take. Anything else is
# dropped rather than emitted.
_REF_ISSUE_SHAPE = re.compile(r"^#[0-9]{1,9}$")
_REF_SHA_SHAPE = re.compile(r"^[0-9a-f]{7,40}$")
# A long undelimited hex run with no declaring context is treated as credential
# material wherever it appears, never as an identifier.
_BARE_SECRET_SHAPE = re.compile(r"^[0-9a-f]{32,}$")
# An event type reads like an identifier, but a stored one is externally
# influenced: any producer that writes the control-plane ``events`` table
# chooses the string. It reaches ``to_dict`` verbatim, so it is validated here
# rather than trusted because of where it came from.
_CP_EVENT_TYPE_SHAPE = re.compile(r"^[A-Za-z][A-Za-z0-9._:+-]{0,63}$")
# Emitted in place of a value that cannot be proven safe. Deliberately not a
# plausible workflow type: an unsafe value is refused, never quietly rewritten
# into a different valid-looking one that would misdescribe the record.
UNSAFE_EVENT_TYPE = "unsafe:redacted"
# Emitted for a CTH heading that is not a declared member of ``CTH_TYPES``. The
# contract is enforced on write (``format_cth_body``) and on assess; the read
# path the timeline uses enforces it too rather than assuming it was.
UNKNOWN_HANDOFF_EVENT_TYPE = "handoff:unrecognized"
# A source record id is a plain integer in both sources it comes from: the
# control-plane ``events`` primary key and a Gitea comment id. ``event_key`` is
# serialized verbatim and is the pagination tiebreak, so anything else is
# refused rather than interpolated into it.
_RECORD_ID_SHAPE = re.compile(r"^[0-9]{1,19}$")
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, ...]:
"""Extract issue/PR and declared-commit references from **redacted** text.
Callers must pass text that has already been through :func:`_redact`; this
function derives a structured field from its input, so extracting ahead of
redaction would republish whatever redaction was about to remove. Every
reference is revalidated by :func:`_validated_evidence_refs` before it is
serialized.
"""
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(1)
if token not in refs:
refs.append(token)
return tuple(refs)
def _validated_evidence_refs(refs: Iterable[str]) -> tuple[tuple[str, ...], bool]:
"""Independently revalidate references immediately before serialization.
Extraction is not trusted on its own. A reference survives only when it has
a known reference shape and is unchanged by a second redaction pass — a
value the redaction policy would alter is credential material that must not
be emitted as a structured field. A full 40-character SHA stays usable
because extraction only accepts a hex run the source text explicitly
declared as a commit. Returns ``(safe_refs, dropped_any)``; ``dropped_any``
marks the event sensitive so the drop is visible rather than silent.
"""
safe: list[str] = []
dropped = False
for ref in refs or ():
try:
token = str(ref).strip()
if not token:
continue
recognised = bool(_REF_ISSUE_SHAPE.match(token) or _REF_SHA_SHAPE.match(token))
if not recognised:
dropped = True
continue
if _redact(token) != token:
dropped = True
continue
if token not in safe:
safe.append(token)
except Exception:
# Fail closed: a reference that cannot be proven safe is dropped.
dropped = True
continue
return (tuple(safe), dropped)
def _safe_session_id(value: Any) -> str | None:
"""Return a session identifier only when it is safe to emit.
The value is authoritative source data — a session the record names for
itself — but it is still free text. It is dropped when redaction alters it
or when it is a bare secret-shaped hex run, so a credential parked in a
session field can never reach the payload or be echoed back by a filter.
"""
if value is None:
return None
text = str(value).strip()
if not text:
return None
if _BARE_SECRET_SHAPE.match(text):
return None
return text if _redact(text) == text else None
def _safe_record_id(value: Any) -> str | None:
"""Return a source record id only when it is a plain numeric identifier.
``event_key`` is serialized verbatim and is the deterministic pagination
tiebreak, so an id is interpolated into it only when it has the shape both
real sources actually produce. A record whose identity cannot be trusted is
refused by the caller rather than keyed on.
"""
if value is None or isinstance(value, bool):
return None
if isinstance(value, int):
return str(value)
text = str(value).strip()
return text if _RECORD_ID_SHAPE.match(text) else None
def _safe_cp_event_type(value: Any) -> tuple[str, bool]:
"""Validate a stored control-plane event type. Returns ``(type, unsafe)``.
The stored value is externally influenced — whichever producer wrote the
``events`` row chose the string — and ``to_dict`` serializes it verbatim, so
it passes a boundary of its own instead of relying on the one ``message``
passes. A value survives only when it is an ordinary identifier, is not a
bare secret-shaped hex run, and is unchanged by a redaction pass. Anything
else fails closed to :data:`UNSAFE_EVENT_TYPE`: the record stays visible as
an audit entry, but the value itself is never republished — not verbatim,
not partially sanitized, and not rewritten into some other valid-looking
type that would misdescribe what happened.
"""
text = ("" if value is None else str(value)).strip()
if not text:
return ("", False)
if _BARE_SECRET_SHAPE.match(text):
return (UNSAFE_EVENT_TYPE, True)
if not _CP_EVENT_TYPE_SHAPE.match(text):
return (UNSAFE_EVENT_TYPE, True)
if _redact(text) != text:
return (UNSAFE_EVENT_TYPE, True)
return (text, False)
def _safe_echo(value: Any) -> Any:
"""Guard a scalar that is echoed back rather than derived from a record.
Query scope and filter values are caller-supplied and are reflected in the
response so an operator can see what was asked. Reflection is still
emission: a value redaction would alter, or a bare secret-shaped hex run, is
replaced by the placeholder instead of being echoed verbatim. Ordinary
scope and filter values pass through untouched.
"""
if value is None or isinstance(value, (int, bool)):
return value
text = str(value)
if _BARE_SECRET_SHAPE.match(text.strip()):
return console_redaction.REDACTED
return _redact(text)
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 = _safe_record_id(row.get("event_id"))
raw_event_type = (row.get("event_type") or "").strip()
if event_id is None or not raw_event_type:
continue
# The stored type is source data, not a trusted constant: validate
# it before it is serialized, exactly as `message` below is redacted
# before it is serialized.
event_type, event_type_unsafe = _safe_cp_event_type(raw_event_type)
kind = row.get("kind")
number = row.get("number")
issue_no, pr_no = _kind_to_numbers(kind, number)
sensitive = event_type_unsafe or any(
hint in raw_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,
# No session_id: the control-plane events table is
# (event_id, work_item_id, event_type, message, created_at)
# and records no session. Inventing one from the work item
# or the message text would be a guess, so this source
# declares the session dimension unsupported instead
# (_SOURCE_FILTER_SUPPORT) and the query layer refuses a
# session filter it cannot honestly answer.
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 is_known_cth_type, parse_cth_comment
# ``kind``/``number`` are interpolated into event_key and correlation_id, so
# they are normalised once here. A scope this adapter cannot express is
# refused outright rather than serialized into an identifier.
kind = (kind or "").strip().lower()
if kind not in ("issue", "pr"):
return []
try:
number = int(number)
except (TypeError, ValueError):
return []
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 ""
comment_id = _safe_record_id(comment.get("id"))
if comment_id is None:
continue
# The CTH heading is free text: the parser accepts whatever follows
# "## CTH:", and only the write and assess paths check it against
# the contract. Check it here too — an unrecognised heading is
# reported as such rather than serialized into event_type, so
# arbitrary, malformed, or secret-shaped heading content has no way
# through. Declared types are preserved exactly.
cth_type_known = is_known_cth_type(cth_type)
# Redaction runs first, and every derived value is taken from the
# redacted text — deriving evidence refs from the raw proof would
# re-emit exactly what redaction was about to remove.
decision = _redact(fields.get("decision"))
proof = _redact(fields.get("proof"))
next_action = _redact(fields.get("next action"))
refs, refs_dropped = _validated_evidence_refs(
_extract_evidence_refs(proof, decision)
)
events.append(
WorkflowEvent(
source=SOURCE_GITEA_HANDOFF,
event_type=(
f"handoff:{cth_type.strip()}"
if cth_type_known
else UNKNOWN_HANDOFF_EVENT_TYPE
),
event_key=f"cth:{kind}:{number}:{comment_id}",
timestamp=_parse_ts(comment.get("created_at")),
actor=_redact((comment.get("user") or {}).get("login")),
role=_redact(fields.get("next owner")),
issue_number=issue_no,
pr_number=pr_no,
# A CTH names its own session when the producer records one;
# it is read from that declared field, never inferred from
# unrelated text.
session_id=_safe_session_id(fields.get("session")),
decision=decision,
message=next_action or _redact(fields.get("status")),
correlation_id=correlation,
evidence_refs=refs,
sensitive=refs_dropped or not cth_type_known,
)
)
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.
``supported_filters`` states which filter dimensions this source's records
can carry; ``unsupported_filters`` names the requested dimensions it cannot,
so an operator can see *why* a source contributed nothing rather than being
left to read an empty list as an absence of activity.
"""
name: str
ok: bool
reason: str | None = None
count: int = 0
supported_filters: tuple[str, ...] = ()
unsupported_filters: tuple[str, ...] = ()
def to_dict(self) -> dict[str, Any]:
return {
"name": self.name,
"ok": self.ok,
"reason": self.reason,
"count": self.count,
"supported_filters": list(self.supported_filters),
"unsupported_filters": list(self.unsupported_filters),
}
def _cp_status(*, ok: bool, reason: str | None = None, count: int = 0) -> SourceStatus:
return SourceStatus(
SOURCE_CONTROL_PLANE,
ok=ok,
# A failure reason is serialized like any other field and is often an
# exception string carrying a path or a transport error, so it crosses
# the redaction boundary too. Static reasons pass through unchanged.
reason=_redact(reason),
count=count,
supported_filters=_SOURCE_FILTER_SUPPORT[SOURCE_CONTROL_PLANE],
)
def _handoff_status(*, ok: bool, reason: str | None = None, count: int = 0) -> SourceStatus:
return SourceStatus(
SOURCE_GITEA_HANDOFF,
ok=ok,
# Same boundary as the control-plane status: this reason can quote an
# error raised by a live authenticated fetch.
reason=_redact(reason),
count=count,
supported_filters=_SOURCE_FILTER_SUPPORT[SOURCE_GITEA_HANDOFF],
)
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 ([], _cp_status(ok=False, reason=f"control-plane DB unavailable: {exc}"))
except sqlite3.Error as exc:
return ([], _cp_status(ok=False, reason=f"control-plane read failed: {exc}"))
finally:
if conn is not None:
conn.close()
events = adapt_cp_events(rows)
return (events, _cp_status(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:
"""One answered timeline query.
``ok`` is False when the query could not be answered as asked — currently
when a requested filter dimension no surviving source can carry was
supplied. The page is then empty *and* the snapshot says so, because an
``ok`` empty page is a claim that no such activity exists.
"""
schema_version: int
remote: str
org: str
repo: str
filters: dict[str, Any]
page: TimelinePage
sources: tuple[SourceStatus, ...]
ok: bool = True
error: dict[str, Any] | None = None
def to_dict(self) -> dict[str, Any]:
return {
"ok": self.ok,
"error": self.error,
"schema_version": self.schema_version,
# Scope and filters are echoed caller input, not derived record
# data. Reflecting a value is still emitting it, so both cross the
# same boundary; ordinary scope and filter values are unchanged.
"scope": {
"remote": _safe_echo(self.remote),
"org": _safe_echo(self.org),
"repo": _safe_echo(self.repo),
},
"filters": {key: _safe_echo(value) for key, value in self.filters.items()},
"sources": [s.to_dict() for s in self.sources],
**self.page.to_dict(),
}
def _unanswerable_reasons(
statuses: Iterable[SourceStatus], unanswerable: Iterable[str]
) -> list[dict[str, str]]:
"""Explain, per source, why each unanswerable dimension went unanswered."""
out: list[dict[str, str]] = []
for status in statuses:
for dim in unanswerable:
if dim not in status.supported_filters:
reason = _SOURCE_FILTER_LIMITS.get(
(status.name, dim), f"this source's records carry no {dim} identity"
)
elif not status.ok:
reason = (
f"this source can carry {dim} but did not run: "
f"{status.reason or 'unavailable'}"
)
else:
continue
out.append({"source": status.name, "filter": dim, "reason": reason})
return out
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.
A filter dimension that no surviving source can carry — a ``session``
filter when the only source that ran is the control plane, whose events
record no session — is refused with ``ok=False`` and a structured error
instead of being answered with an empty page.
"""
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(
_handoff_status(
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(
_handoff_status(
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(_handoff_status(ok=True, count=len(handoff_events)))
except Exception as exc: # fail soft: a fetch/parse error degrades this source only
statuses.append(_handoff_status(ok=False, reason=f"handoff source failed: {exc}"))
requested = tuple(
name
for name, value in ((FILTER_ISSUE, issue), (FILTER_PR, pr), (FILTER_SESSION, session))
if value is not None
)
statuses = [
replace(
status,
unsupported_filters=tuple(
dim for dim in requested if dim not in status.supported_filters
),
)
for status in statuses
]
filters = {"issue": issue, "pr": pr, "session": session}
# A dimension is answerable only if a source that actually ran can carry it.
# If none can, refuse: an empty page would assert "no such activity", which
# is a claim this timeline is not in a position to make.
answerable: set[str] = set()
for status in statuses:
if status.ok:
answerable.update(status.supported_filters)
unanswerable = tuple(dim for dim in requested if dim not in answerable)
if unanswerable:
return TimelineSnapshot(
schema_version=TIMELINE_SCHEMA_VERSION,
remote=remote,
org=org,
repo=repo,
filters=filters,
page=paginate([], limit=limit, offset=offset),
sources=tuple(statuses),
ok=False,
error={
"code": "filter_not_supported",
"unsupported_filters": list(unanswerable),
"detail": (
"no timeline source that ran can answer "
+ ", ".join(f"'{dim}'" for dim in unanswerable)
+ "; the result is refused rather than returned empty"
),
"sources": _unanswerable_reasons(statuses, unanswerable),
},
)
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=filters,
page=page,
sources=tuple(statuses),
)
def snapshot_to_dict(snapshot: TimelineSnapshot) -> dict[str, Any]:
return snapshot.to_dict()