Compare commits

...
Author SHA1 Message Date
sysadmin 38506a753f Merge remote-tracking branch 'prgs/master' into feat/issue-651-usage-cost-analytics
# Conflicts:
#	docs/webui-authz-audit.md
#	webui/console_authz.py
2026-07-24 08:33:35 -04:00
sysadmin 103d0df289 Merge pull request 'feat(webui): sanctioned restart and graceful reload controls (Closes #642)' (#863) from feat/issue-642-sanctioned-restart-controls into master 2026-07-24 07:21:17 -05:00
sysadminandClaude Opus 4.8 b179610e7f fix(webui): remediate PR #876 REQUEST_CHANGES for analytics (#651)
Address review blockers and medium/low findings:

- F1: HTML-escape all dynamic analytics fields (html.escape quote=True)
- F2: Gate POST /api/v1/analytics/usage through console_authz
  record_analytics_usage (fail closed for unauthenticated / Phase 1)
- F3: Enforce usage_events retention (max rows + max age)
- F4: Coerce None remote/org/repo to empty strings in load_analytics

Add tests for XSS escaping, unauthorized ingest deny, retention, and
None scope coercion. Document authz + retention in analytics guide.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-24 08:11:12 -04:00
jcwalker3 205207abb0 Merge branch 'master' into feat/issue-642-sanctioned-restart-controls 2026-07-24 07:03:50 -05:00
jcwalker3 0041542fe2 Merge branch 'master' into feat/issue-651-usage-cost-analytics 2026-07-24 07:03:24 -05:00
sysadmin a87a7d1da2 Merge pull request 'Implement native author issue worktree bootstrap (#850)' (#853) from fix/issue-850-native-mcp-bootstrap into master 2026-07-24 07:02:48 -05:00
jcwalker3 18bca47977 Merge branch 'master' into feat/issue-642-sanctioned-restart-controls 2026-07-24 06:55:15 -05:00
jcwalker3 1f144705e2 Merge branch 'master' into feat/issue-651-usage-cost-analytics 2026-07-24 06:55:03 -05:00
sysadmin 5c5c1fdf77 feat(webui): model usage, token cost, latency, and performance analytics (Closes #651) 2026-07-24 07:16:44 -04:00
jcwalker3 2a5d6571ec Merge branch 'master' into feat/issue-642-sanctioned-restart-controls 2026-07-24 06:16:03 -05:00
jcwalker3 37c3e5dc39 Merge branch 'master' into feat/issue-642-sanctioned-restart-controls 2026-07-24 06:11:00 -05:00
jcwalker3 d2eaca4949 Merge branch 'master' into feat/issue-642-sanctioned-restart-controls 2026-07-24 05:54:59 -05:00
jcwalker3 4ff4d2acc9 Merge branch 'master' into feat/issue-642-sanctioned-restart-controls 2026-07-24 02:07:44 -05:00
jcwalker3 fc8fe329d2 Merge branch 'master' into feat/issue-642-sanctioned-restart-controls 2026-07-23 22:22:36 -05:00
sysadminandClaude Opus 4.8 e593444eea feat(webui): sanctioned restart and graceful reload controls (Closes #642)
Adds a gated `system.restart_namespace` action so operators and workers can
restart or gracefully reload MCP namespaces through an authorized, audited
path instead of manual host process killing (#630).

- webui/sanctioned_restart.py: restart/reload operation model, dry-run
  intent preview, confirmation-string enforcement, audit emission, and
  post-restart health verification. Fails closed on unknown auth, missing
  capability, or ambiguous target namespace.
- webui/console_authz.py: RBAC entries for the restart capability with
  secret redaction preserved.
- webui/gated_actions.py: registers the restart action in the gated action
  framework so it cannot be invoked without capability + confirmation.
- task_capability_map.py: capability mapping for the restart operation.
- docs/sanctioned-restart-controls.md: operator documentation for the
  sanctioned path and the explicit prohibition on pkill recovery.
- docs/webui-authz-audit.md: audit model updated for restart events.
- tests/test_webui_sanctioned_restart.py: authorized preview, unauthorized
  deny, confirmation enforcement, contamination classification, and audit
  emission coverage.

No unrestricted kill path is exposed; manual pkill remains classified as
contamination and continues to block clean claims.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-23 23:19:21 -04:00
15 changed files with 2754 additions and 4 deletions
+231 -1
View File
@@ -31,7 +31,7 @@ from typing import Any, Iterator, Sequence
import dependency_graph
SCHEMA_VERSION = 4
SCHEMA_VERSION = 5
# Assignable work kinds only — raw monitoring incidents are never work items.
WORK_KINDS = frozenset({"issue", "pr"})
@@ -186,6 +186,34 @@ CREATE INDEX IF NOT EXISTS idx_dependency_edges_target
ON dependency_edges(remote, org, repo, target_kind, target_number);
CREATE INDEX IF NOT EXISTS idx_assignments_session ON assignments(session_id, status);
CREATE INDEX IF NOT EXISTS idx_incident_gitea ON incident_links(gitea_org, gitea_repo, gitea_issue_number);
-- Model usage, token cost, latency, and performance events (#651)
CREATE TABLE IF NOT EXISTS usage_events (
usage_id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT,
remote TEXT NOT NULL DEFAULT 'dadeschools',
org TEXT NOT NULL DEFAULT '',
repo TEXT NOT NULL DEFAULT '',
project_id TEXT,
role TEXT NOT NULL DEFAULT 'unknown',
model TEXT NOT NULL DEFAULT 'unknown',
issue_number INTEGER,
pr_number INTEGER,
stage TEXT NOT NULL DEFAULT 'unknown',
input_tokens INTEGER,
output_tokens INTEGER,
total_tokens INTEGER,
estimated_cost_usd REAL,
latency_ms INTEGER,
duration_ms INTEGER,
status TEXT NOT NULL DEFAULT 'success',
metadata TEXT,
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_usage_events_scope ON usage_events(remote, org, repo);
CREATE INDEX IF NOT EXISTS idx_usage_events_role_model ON usage_events(role, model);
CREATE INDEX IF NOT EXISTS idx_usage_events_stage ON usage_events(stage);
"""
@@ -339,6 +367,7 @@ class ControlPlaneDB:
self._migrate_incident_links_null_scope(conn)
self._migrate_lease_lifecycle_columns(conn)
self._migrate_session_ownership_columns(conn)
self._migrate_usage_events_table(conn)
conn.execute(
"INSERT OR REPLACE INTO schema_meta(key, value) VALUES (?, ?)",
("schema_version", str(SCHEMA_VERSION)),
@@ -518,6 +547,207 @@ class ControlPlaneDB:
f"UPDATE incident_links SET {col} = '' WHERE {col} IS NULL"
)
def _migrate_usage_events_table(self, conn: sqlite3.Connection) -> None:
"""Create usage_events table and indexes if they do not exist (#651)."""
conn.execute("""
CREATE TABLE IF NOT EXISTS usage_events (
usage_id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT,
remote TEXT NOT NULL DEFAULT 'dadeschools',
org TEXT NOT NULL DEFAULT '',
repo TEXT NOT NULL DEFAULT '',
project_id TEXT,
role TEXT NOT NULL DEFAULT 'unknown',
model TEXT NOT NULL DEFAULT 'unknown',
issue_number INTEGER,
pr_number INTEGER,
stage TEXT NOT NULL DEFAULT 'unknown',
input_tokens INTEGER,
output_tokens INTEGER,
total_tokens INTEGER,
estimated_cost_usd REAL,
latency_ms INTEGER,
duration_ms INTEGER,
status TEXT NOT NULL DEFAULT 'success',
metadata TEXT,
created_at TEXT NOT NULL
);
""")
conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_scope ON usage_events(remote, org, repo);")
conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_role_model ON usage_events(role, model);")
conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_stage ON usage_events(stage);")
# #651 retention: cap growth so unauthenticated or high-volume ingest
# cannot DoS the control-plane DB (PR #876 F3). Applied after every write.
USAGE_EVENTS_MAX_ROWS = 10_000
USAGE_EVENTS_RETENTION_DAYS = 90
def record_usage_event(
self,
*,
session_id: str | None = None,
remote: str = "dadeschools",
org: str = "",
repo: str = "",
project_id: str | None = None,
role: str = "unknown",
model: str = "unknown",
issue_number: int | None = None,
pr_number: int | None = None,
stage: str = "unknown",
input_tokens: int | None = None,
output_tokens: int | None = None,
total_tokens: int | None = None,
estimated_cost_usd: float | None = None,
latency_ms: int | None = None,
duration_ms: int | None = None,
status: str = "success",
metadata: str | dict[str, Any] | None = None,
created_at: str | None = None,
) -> int:
"""Record a model usage, token cost, latency, or stage performance event (#651)."""
ts = created_at or _ts()
meta_str: str | None = None
if metadata is not None:
from webui import console_redaction
redacted_meta = console_redaction.redact_payload(metadata)
if isinstance(redacted_meta, str):
meta_str = redacted_meta
else:
try:
meta_str = json.dumps(redacted_meta, default=str)
except Exception:
meta_str = str(redacted_meta)
if total_tokens is None and (input_tokens is not None or output_tokens is not None):
total_tokens = (input_tokens or 0) + (output_tokens or 0)
with self._tx(immediate=True) as conn:
cursor = conn.execute(
"""
INSERT INTO usage_events (
session_id, remote, org, repo, project_id, role, model,
issue_number, pr_number, stage, input_tokens, output_tokens,
total_tokens, estimated_cost_usd, latency_ms, duration_ms,
status, metadata, created_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
session_id,
remote,
org,
repo,
project_id,
role,
model,
issue_number,
pr_number,
stage,
input_tokens,
output_tokens,
total_tokens,
estimated_cost_usd,
latency_ms,
duration_ms,
status,
meta_str,
ts,
),
)
usage_id = cursor.lastrowid
self._enforce_usage_events_retention(conn)
return usage_id
def _enforce_usage_events_retention(self, conn: sqlite3.Connection) -> None:
"""Drop aged and excess usage_events rows (PR #876 F3)."""
# Age-based: ISO-8601 UTC timestamps compare lexicographically.
cutoff = (
datetime.now(timezone.utc)
- timedelta(days=int(self.USAGE_EVENTS_RETENTION_DAYS))
).strftime("%Y-%m-%dT%H:%M:%SZ")
conn.execute(
"DELETE FROM usage_events WHERE created_at < ?",
(cutoff,),
)
# Count-based: keep the newest USAGE_EVENTS_MAX_ROWS by usage_id.
max_rows = int(self.USAGE_EVENTS_MAX_ROWS)
if max_rows > 0:
conn.execute(
"""
DELETE FROM usage_events
WHERE usage_id NOT IN (
SELECT usage_id FROM usage_events
ORDER BY usage_id DESC
LIMIT ?
)
""",
(max_rows,),
)
def query_usage_events(
self,
*,
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
project_id: str | None = None,
role: str | None = None,
model: str | None = None,
issue_number: int | None = None,
pr_number: int | None = None,
stage: str | None = None,
session_id: str | None = None,
limit: int = 500,
offset: int = 0,
) -> list[dict[str, Any]]:
"""Query stored usage events matching filters (#651)."""
conditions = []
params = []
if remote:
conditions.append("remote = ?")
params.append(remote)
if org:
conditions.append("org = ?")
params.append(org)
if repo:
conditions.append("repo = ?")
params.append(repo)
if project_id:
conditions.append("project_id = ?")
params.append(project_id)
if role:
conditions.append("role = ?")
params.append(role)
if model:
conditions.append("model = ?")
params.append(model)
if issue_number is not None:
conditions.append("issue_number = ?")
params.append(issue_number)
if pr_number is not None:
conditions.append("pr_number = ?")
params.append(pr_number)
if stage:
conditions.append("stage = ?")
params.append(stage)
if session_id:
conditions.append("session_id = ?")
params.append(session_id)
where_clause = f"WHERE {' AND '.join(conditions)}" if conditions else ""
sql = f"""
SELECT * FROM usage_events
{where_clause}
ORDER BY usage_id ASC
LIMIT ? OFFSET ?
"""
params.extend([limit, offset])
with self._tx(immediate=False) as conn:
cursor = conn.execute(sql, params)
rows = cursor.fetchall()
return [dict(row) for row in rows]
# ── sessions ──────────────────────────────────────────────────────────
def upsert_session(
@@ -0,0 +1,143 @@
# Model Usage, Token Cost, Latency, and Workflow Analytics (Phase 4)
- **Tracking Issue:** [#651](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/651)
- **Parent Epic:** [#631](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/651)
- **Console Surface:** `/analytics`, `/api/v1/analytics`, `/api/v1/analytics/usage`
## 1. Overview
The Web Console Analytics module provides durable, aggregate visibility into **model usage, token cost, latency percentiles, and workflow-stage performance** across projects, worker roles, AI models, issues, and PRs.
### Non-Goals
- No mandatory client-side telemetry that leaks prompts or secret keys.
- No third-party payment provider or billing integration.
- No automatic model routing changes without controller policy (#647).
---
## 2. Event Schema (`usage_events`)
Usage metrics are stored in the control-plane database under table `usage_events`.
| Column | Type | Description |
|---|---|---|
| `usage_id` | `INTEGER` | Primary key (autoincrement) |
| `session_id` | `TEXT` | Optional active session identifier |
| `remote` | `TEXT` | Known Gitea instance (`dadeschools` or `prgs`) |
| `org` | `TEXT` | Repository owner / organization |
| `repo` | `TEXT` | Repository name |
| `project_id` | `TEXT` | Optional project identifier |
| `role` | `TEXT` | Active worker role (`author`, `reviewer`, `merger`, `reconciler`, `controller`) |
| `model` | `TEXT` | LLM model identifier (e.g. `gemini-3.6-flash`, `claude-3-5-sonnet`) |
| `issue_number` | `INTEGER` | Correlated Gitea issue number (optional) |
| `pr_number` | `INTEGER` | Correlated Gitea PR number (optional) |
| `stage` | `TEXT` | Workflow stage (`preflight`, `implementation`, `review`, `merge`, `reconciliation`) |
| `input_tokens` | `INTEGER` | Input token count (optional / nullable) |
| `output_tokens` | `INTEGER` | Output token count (optional / nullable) |
| `total_tokens` | `INTEGER` | Total token count (optional / nullable) |
| `estimated_cost_usd` | `REAL` | Estimated USD cost (optional / nullable) |
| `latency_ms` | `INTEGER` | Request latency in milliseconds (optional / nullable) |
| `duration_ms` | `INTEGER` | Stage execution duration in milliseconds (optional / nullable) |
| `status` | `TEXT` | Outcome status (`success`, `failure`, `timeout`) |
| `metadata` | `TEXT` | Redacted metadata or summary string |
| `created_at` | `TEXT` | ISO 8601 UTC timestamp |
---
## 3. Handling of Missing Data ("Unknown" vs. Zero Fabrication)
To ensure operational metrics accurately reflect evidence:
- **Untracked or missing metrics are displayed as `Unknown`**, never zero-fabricated.
- If an event omits `estimated_cost_usd`, `latency_ms`, or token counts, the aggregator marks those fields as missing (`None`) rather than defaulting to `0` or `$0.00`.
- Summary tables and KPI cards explicitly indicate when data is unmeasured or partially reported.
---
## 4. Redaction & Security Rules
Per `#633` security policy:
- Free-text fields (`metadata`, `prompt_summary`, `session_id`) are run through `console_redaction.redact_text` before persistence and output serialization.
- Secret tokens, keychain commands, authorization headers, passwords, and JWTs are stripped automatically.
---
## 5. Opt-in Instrumentation Guide
Applications, MCP servers, and background sessions can report usage metrics through either Python API or HTTP ingestion.
### Python Ingestion
```python
from webui.analytics_loader import record_usage
record_usage(
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
role="author",
model="gemini-3.6-flash",
issue_number=651,
stage="implementation",
input_tokens=1420,
output_tokens=380,
total_tokens=1800,
estimated_cost_usd=0.00045,
latency_ms=320,
duration_ms=4500,
status="success",
metadata={"note": "Implementation of analytics module"},
)
```
### HTTP Ingestion API (authorized write)
`POST /api/v1/analytics/usage` is a **gated write**. It runs through
`console_authz` action `record_analytics_usage` (operator+, Phase 2 execution).
Unauthenticated or phase-inactive requests receive **403** and do not write.
Prefer in-process `record_usage` for MCP / session instrumentation.
```http
POST /api/v1/analytics/usage HTTP/1.1
Content-Type: application/json
# Requires authenticated principal with record_analytics_usage execution enabled
{
"remote": "dadeschools",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
"role": "author",
"model": "gemini-3.6-flash",
"issue_number": 651,
"stage": "implementation",
"input_tokens": 1420,
"output_tokens": 380,
"total_tokens": 1800,
"estimated_cost_usd": 0.00045,
"latency_ms": 320,
"duration_ms": 4500,
"status": "success",
"metadata": "Analytics schema landed"
}
```
### Retention
`usage_events` is retained with hard caps applied on every write:
| Limit | Default |
|---|---|
| Max rows | 10,000 (`ControlPlaneDB.USAGE_EVENTS_MAX_ROWS`) |
| Max age | 90 days (`ControlPlaneDB.USAGE_EVENTS_RETENTION_DAYS`) |
Older rows (by `created_at`) and excess oldest rows (by `usage_id`) are deleted
after each insert so unbounded growth / DoS-by-volume cannot fill the DB.
---
## 6. Querying Analytics API
```http
GET /api/v1/analytics?role=author&stage=implementation HTTP/1.1
```
Returns `AnalyticsSnapshot` JSON containing aggregations (`by_model`, `by_stage`, `by_role`, `by_work_item`, `by_project`) and latency percentiles (`p50`, `p90`, `p95`, `p99`).
+122
View File
@@ -0,0 +1,122 @@
# Sanctioned restart and graceful reload controls (#642)
Sessions used to recover MCP connectivity by killing the host daemon
(`pkill -f mcp_server.py`, #630). That is forbidden and stays forbidden: it
kills every namespace on the host, contaminates whichever session survives, and
leaves no audit trail. This document describes the sanctioned replacement,
implemented in `webui/sanctioned_restart.py`.
## What the console will and will not do
The console **never** restarts anything. It authorizes an intent, records it,
and hands off to a host supervisor. There is no code path in which the console
sends a signal, spawns a process, or renders a kill command — a regression test
asserts the module contains no `subprocess`, `signal`, `os.kill`, `os.system`,
or `popen` reference, and that no returned payload contains a kill command.
## Operations
| Mode | Action | Minimum role | Behaviour |
|------|--------|--------------|-----------|
| `reload` | `system.reload_namespace` | controller | Host supervisor reloads the namespace in place, draining in-flight requests. |
| `restart` | `system.restart_namespace` | admin | Host supervisor restarts the namespace. In-flight requests are lost. |
Scope is always exactly one namespace. A fleet-wide restart is an explicit
non-goal: `all`, `*`, `fleet`, and an empty scope are refused with
`fleet_scope_not_permitted`, because that is precisely the blast radius the
forbidden kill already had. An unrecognised namespace is refused rather than
passed through to the host.
## The gate sequence
`assess_restart_request()` applies every gate in order and reports the first
failure with a stable reason code:
| Order | Gate | Reason code on failure |
|-------|------|------------------------|
| 1 | Mode is `restart` or `reload` | `unknown_mode` |
| 2 | Scope is a single known namespace | `fleet_scope_not_permitted`, `unknown_namespace` |
| 3 | Principal holds the required console role | `unauthorized` |
| 4 | Confirmation phrase supplied | `confirmation_required` |
| 5 | Confirmation names this namespace and mode | `confirmation_mismatch` |
| 6 | Out-of-band operator authorization present | `operator_authorization_missing` |
| 7 | Runtime is not contaminated | `contaminated_runtime` |
| 8 | Host restart hook configured | `restart_hook_not_configured` |
Passing every gate yields `host_action_required`, never "restarted".
### Confirmation binds the namespace
The required phrase is `"<mode> <namespace>"` — for example
`restart gitea-author`. Binding the namespace into the phrase is the point: a
confirmation typed for one namespace cannot be replayed against another.
### Operator authorization is not self-assertable
Host daemon maintenance is authorized out of band through
`GITEA_OPERATOR_DAEMON_MAINTENANCE_AUTHORIZATION`, read from the process
environment and nowhere else (#630; #710 finding F1). A worker session cannot
set an environment variable for an already-running daemon, so this cannot be
faked the way a tool argument could.
### The host hook
`GITEA_SANCTIONED_RESTART_HOOK` holds an opaque reference the *host* resolves —
a supervisor label such as a launchd job name, never a command line. With no
hook configured the request is refused; the console does not fall back to a
process kill. The value is read server-side and never rendered to a client.
## Manual kill remains contamination
`classify_restart_command()` classifies an operator-proposed recovery command.
A manual `pkill`/`kill`/`killall` of the MCP daemon is contamination, not a
restart: it returns `clean_claim_allowed: false` and builds a durable
contamination marker (redacted command only, never secrets) naming
`system.restart_namespace` as the sanctioned alternative.
A live, uncleared contamination marker also blocks a restart. This is stricter
than #630's task-scoped gate, which deliberately lets a contaminated worker keep
commenting and handing off: restarting a contaminated runtime would launder the
contamination rather than resolve it. Clear the marker through the reconciler
path first.
## Post-restart health verification
After the host supervisor acts, `verify_post_restart_health()` decides whether
the session may claim to be clean:
| Status | Meaning | Clean claim |
|--------|---------|-------------|
| `clean` | Required tool callable, proven through the live client namespace | Allowed |
| `unproven` | Reported healthy without live client-namespace evidence | Refused |
| `unhealthy` | Probe failed | Refused |
Only `probe_source=client_namespace` evidence clears a session. Static tool
registration is not proof, and neither is an offline subprocess probe — an IDE
client can hold a registered tool list while live calls fail with
`client is closing: EOF` (see
[`mcp-namespace-health.md`](mcp-namespace-health.md)).
## Audit
Every attempt — allowed or denied — is recorded through
`webui.console_audit` with actor, target namespace, mode, result, and reason
code, and is redacted before it is persisted. `system.restart_namespace` is
break-glass, so its records are retained for 730 days. Records carry
`process_kill_executed: false`, which is a fact about the code path rather than
a claim: no such path exists.
## Environment variables
| Variable | Purpose |
|----------|---------|
| `GITEA_SANCTIONED_RESTART_HOOK` | Host supervisor reference; absent means restart is refused. |
| `GITEA_OPERATOR_DAEMON_MAINTENANCE_AUTHORIZATION` | Out-of-band operator authorization reference. |
| `WEBUI_AUDIT_LOG` | Console audit sink; absent means records are built but not persisted. |
## Non-goals
* No unrestricted `kill` from the UI, in any role, in any phase.
* No fleet-wide restart.
* No silent auto-restart loop: every attempt is confirmed and audited.
* This does not implement the Phase 1 health API (#634).
+12 -2
View File
@@ -91,6 +91,9 @@ already define, and a regression test asserts each mapping matches.
| `close_pr` | controller | privileged | `gitea.pr.close` | Yes | No | No | 3 |
| `merge_pr` | controller | privileged | `gitea.pr.merge` | Yes | **Yes** | **Yes** | 3 |
| `delete_branch` | admin | destructive | `gitea.branch.delete` | Yes | **Yes** | **Yes** | 3 |
| `record_analytics_usage` | operator | gated_write | `runtime.record_analytics_usage` | Yes | No | No | 2 |
| `system.reload_namespace` | controller | privileged | `runtime.reload_namespace` | Yes | No | No | 2 |
| `system.restart_namespace` | admin | destructive | `runtime.restart_namespace` | Yes | **Yes** | **Yes** | 2 |
**Dual control** means the acting principal may not be the sole authority: a
second distinct principal must confirm. **Break-glass** means the action is
@@ -102,6 +105,13 @@ honouring it.
`delete_branch` is admin-only rather than controller because it is the one
irreversible action in the set.
`system.restart_namespace` is admin-only for the same reason: restarting a
namespace drops every in-flight request on it. `system.reload_namespace` drains
first, so it is privileged but not destructive. Neither action is ever executed
by the console — both hand off to a host supervisor, and neither exposes a raw
process kill. See
[`sanctioned-restart-controls.md`](sanctioned-restart-controls.md) (#642).
### Authorization decision
`authorize(action_id, principal, for_execution=False)` returns a decision
@@ -199,8 +209,8 @@ breaking the request it describes.
| Class | Applies to | Default |
|-------|-----------|---------|
| `standard` | Routine gated writes | 90 days |
| `privileged` | `review_pr`, `close_pr`, and any unclassifiable action | 365 days |
| `break_glass` | `merge_pr`, `delete_branch` | 730 days |
| `privileged` | `review_pr`, `close_pr`, `system.reload_namespace`, and any unclassifiable action | 365 days |
| `break_glass` | `merge_pr`, `delete_branch`, `system.restart_namespace` | 730 days |
Each record carries its own class, day count, and computed `expires_at`, so
retention is auditable per record rather than inferred from file age. An
+24
View File
@@ -363,6 +363,21 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"role": "controller",
},
# #642: sanctioned host-daemon lifecycle controls. Deliberately *not* a
# ``gitea.*`` operation — restarting an MCP namespace is a host action, not
# a Gitea API call, and no configured Gitea profile should be able to
# satisfy it by accident. Authority comes from the console RBAC model plus
# out-of-band operator authorization (#630); these entries exist so the
# console cannot invent an authority the capability layer never declared.
"restart_namespace": {
"permission": "runtime.restart_namespace",
"role": "controller",
},
"reload_namespace": {
"permission": "runtime.reload_namespace",
"role": "controller",
},
# #601 first-class lease lifecycle — inspect/list need read; mutations gate on
# ownership in the control-plane DB (not a separate Gitea write permission).
"list_workflow_leases": {
@@ -495,6 +510,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.issue.comment",
"role": "author",
},
# #651 console analytics ingest — control-plane DB write, not a Gitea API
# call. Authority comes from console RBAC (operator+) plus phase gating;
# permission string is a non-Gitea runtime capability so no Gitea profile
# can satisfy it by accident.
"record_analytics_usage": {
"permission": "runtime.record_analytics_usage",
"role": "author",
},
}
+1 -1
View File
@@ -36,7 +36,7 @@ class ControlPlaneDBTest(unittest.TestCase):
rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall())
finally:
conn.close()
self.assertEqual(rows["schema_version"], "4")
self.assertEqual(rows["schema_version"], "5")
self.assertIn("DB coordinates", rows["architecture"])
self.assertIn("bridge", rows["architecture"].lower())
+252
View File
@@ -0,0 +1,252 @@
"""Unit and integration tests for Model Usage & Performance Analytics (#651)."""
from __future__ import annotations
import os
import tempfile
import unittest
from starlette.testclient import TestClient
import control_plane_db
from webui.analytics_loader import (
ANALYTICS_SCHEMA_VERSION,
compute_percentile,
load_analytics,
record_usage,
)
from webui.app import create_app
from webui import console_redaction
class AnalyticsLoaderTest(unittest.TestCase):
def setUp(self) -> None:
self.temp_dir = tempfile.TemporaryDirectory()
self.db_path = os.path.join(self.temp_dir.name, "test_control_plane.sqlite3")
os.environ["GITEA_CONTROL_PLANE_DB"] = self.db_path
self.db = control_plane_db.ControlPlaneDB(db_path=self.db_path)
def tearDown(self) -> None:
self.temp_dir.cleanup()
def test_compute_percentile(self) -> None:
self.assertIsNone(compute_percentile([], 50.0))
self.assertEqual(compute_percentile([100], 50.0), 100.0)
# 2 elements: [100, 200]
self.assertEqual(compute_percentile([100, 200], 50.0), 150.0)
# 100 elements: 1..100
vals = list(range(1, 101))
self.assertEqual(compute_percentile(vals, 50.0), 50.5)
self.assertAlmostEqual(compute_percentile(vals, 90.0), 90.1)
def test_record_and_aggregate_usage(self) -> None:
# Record event 1 (complete data)
u1 = record_usage(
db_path=self.db_path,
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
role="author",
model="gemini-3.6-flash",
issue_number=651,
stage="implementation",
input_tokens=1000,
output_tokens=500,
estimated_cost_usd=0.0015,
latency_ms=200,
duration_ms=3000,
metadata={"secret_key": "secret123", "note": "token=secret123"},
)
self.assertGreater(u1, 0)
# Record event 2 (missing tokens and cost -> unknown)
u2 = record_usage(
db_path=self.db_path,
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
role="reviewer",
model="claude-3-5-sonnet",
pr_number=846,
stage="review",
latency_ms=500,
duration_ms=6000,
)
self.assertGreater(u2, u1)
snapshot = load_analytics(
db_path=self.db_path,
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
)
self.assertTrue(snapshot.ok)
self.assertEqual(snapshot.schema_version, ANALYTICS_SCHEMA_VERSION)
self.assertEqual(snapshot.total_events, 2)
# Verify overall summary
summary = snapshot.overall_summary
self.assertEqual(summary.total_events, 2)
self.assertEqual(summary.events_with_tokens, 1)
self.assertEqual(summary.total_tokens, 1500)
self.assertEqual(summary.events_with_cost, 1)
self.assertEqual(summary.estimated_cost_usd, 0.0015)
self.assertEqual(summary.events_with_latency, 2)
self.assertEqual(summary.latency_p50_ms, 350.0)
# Verify missing data handling (AC 3: not zero-fabricated)
reviewer_model = snapshot.by_model.get("claude-3-5-sonnet")
self.assertIsNotNone(reviewer_model)
self.assertEqual(reviewer_model.total_events, 1)
self.assertEqual(reviewer_model.events_with_tokens, 0)
self.assertIsNone(reviewer_model.total_tokens)
self.assertEqual(reviewer_model.display_tokens, "Unknown")
self.assertEqual(reviewer_model.events_with_cost, 0)
self.assertIsNone(reviewer_model.estimated_cost_usd)
self.assertEqual(reviewer_model.display_cost, "Unknown")
# Verify redaction (AC 4)
e1 = [e for e in snapshot.events if e.usage_id == u1][0]
self.assertIsNotNone(e1.metadata)
self.assertNotIn("secret123", e1.metadata)
self.assertIn("[REDACTED]", e1.metadata)
def test_missing_db_fail_soft(self) -> None:
invalid_path = "/nonexistent_path_dir/db.sqlite3"
snapshot = load_analytics(db_path=invalid_path)
self.assertFalse(snapshot.ok)
self.assertIn("control_plane_db_unavailable", snapshot.reason)
self.assertEqual(snapshot.overall_summary.display_tokens, "Unknown")
class AnalyticsWebUITest(unittest.TestCase):
def setUp(self) -> None:
self.temp_dir = tempfile.TemporaryDirectory()
self.db_path = os.path.join(self.temp_dir.name, "test_webui.sqlite3")
os.environ["GITEA_CONTROL_PLANE_DB"] = self.db_path
self.app = create_app()
self.client = TestClient(self.app)
record_usage(
db_path=self.db_path,
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
role="author",
model="gemini-3.6-flash",
issue_number=651,
stage="implementation",
input_tokens=2000,
output_tokens=1000,
estimated_cost_usd=0.003,
latency_ms=150,
duration_ms=2500,
)
def tearDown(self) -> None:
self.temp_dir.cleanup()
def test_analytics_html_route(self) -> None:
response = self.client.get("/analytics")
self.assertEqual(response.status_code, 200)
self.assertIn("Model Usage & Performance Analytics", response.text)
self.assertIn("gemini-3.6-flash", response.text)
self.assertIn("3,000", response.text)
def test_analytics_api_route(self) -> None:
response = self.client.get("/api/v1/analytics")
self.assertEqual(response.status_code, 200)
data = response.json()
self.assertTrue(data["ok"])
self.assertEqual(data["total_events"], 1)
self.assertIn("gemini-3.6-flash", data["by_model"])
def test_analytics_ingest_unauthorized_denied(self) -> None:
"""F2: unauthenticated POST must not write the control-plane DB."""
payload = {
"remote": "dadeschools",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
"role": "reviewer",
"model": "claude-3-5-sonnet",
"pr_number": 846,
"stage": "review",
"input_tokens": 500,
"output_tokens": 100,
"latency_ms": 400,
"metadata": "Review note token=secret456",
}
response = self.client.post("/api/v1/analytics/usage", json=payload)
self.assertEqual(response.status_code, 403)
res_json = response.json()
self.assertFalse(res_json.get("ok", True))
self.assertEqual(res_json.get("error"), "unauthorized")
authorization = res_json.get("authorization") or {}
self.assertFalse(authorization.get("allowed"))
self.assertFalse(authorization.get("execution_enabled"))
# No new row written
res2 = self.client.get("/api/v1/analytics")
self.assertEqual(res2.status_code, 200)
self.assertEqual(res2.json()["total_events"], 1)
def test_html_escapes_script_bearing_model_role_stage(self) -> None:
"""F1: stored XSS — dynamic model/role/stage must render escaped."""
xss = '<script>alert(1)</script>'
record_usage(
db_path=self.db_path,
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
role=xss,
model=xss,
stage=xss,
issue_number=999,
status="success",
)
response = self.client.get("/analytics")
self.assertEqual(response.status_code, 200)
# Raw tag must not appear; escaped form must.
self.assertNotIn("<script>alert(1)</script>", response.text)
self.assertIn("&lt;script&gt;alert(1)&lt;/script&gt;", response.text)
def test_load_analytics_coerces_none_scope(self) -> None:
"""F4: None remote/org/repo become empty strings, never None."""
snapshot = load_analytics(db_path=self.db_path, remote=None, org=None, repo=None)
self.assertIsInstance(snapshot.remote, str)
self.assertIsInstance(snapshot.org, str)
self.assertIsInstance(snapshot.repo, str)
self.assertEqual(snapshot.remote, "")
self.assertEqual(snapshot.org, "")
self.assertEqual(snapshot.repo, "")
def test_usage_events_retention_max_rows(self) -> None:
"""F3: record_usage_event enforces USAGE_EVENTS_MAX_ROWS."""
db = control_plane_db.ControlPlaneDB(db_path=self.db_path)
original_max = db.USAGE_EVENTS_MAX_ROWS
try:
db.USAGE_EVENTS_MAX_ROWS = 3
for i in range(5):
db.record_usage_event(
remote="dadeschools",
org="org",
repo="repo",
role="author",
model=f"model-{i}",
stage="test",
)
rows = db.query_usage_events(limit=100)
self.assertLessEqual(len(rows), 3)
# Newest three retained
models = {r["model"] for r in rows}
self.assertEqual(models, {"model-2", "model-3", "model-4"})
finally:
db.USAGE_EVENTS_MAX_ROWS = original_max
if __name__ == "__main__":
unittest.main()
+502
View File
@@ -0,0 +1,502 @@
"""Sanctioned restart / graceful reload control tests (#642).
Acceptance criteria under test:
1. The sanctioned restart path is implemented behind gates (capability,
confirmation, operator authorization, host hook).
2. Manual ``pkill`` stays forbidden and is classified as contamination.
3. Post-restart mutations require clean health/session proof.
4. Authorized restart preview, unauthorized deny, contamination classification.
5. No entry point exposes a raw kill.
"""
import json
import os
import tempfile
import unittest
import mcp_namespace_health
import runtime_recovery_guard
from task_capability_map import TASK_CAPABILITY_MAP
from webui import console_audit, console_authz, gated_actions, sanctioned_restart
NAMESPACE = "gitea-author"
# An operator-authorized, hook-configured host. Passed explicitly so no test
# depends on (or mutates) the real process environment.
READY_ENV = {
sanctioned_restart.RESTART_HOOK_ENV: "launchd:cc.prgs.gitea-author",
runtime_recovery_guard.OPERATOR_AUTHORIZATION_ENV: "ops-ticket-4821",
}
def admin(subject: str = "[email protected]") -> console_authz.Principal:
return console_authz.Principal(
subject=subject,
role=console_authz.ADMIN,
identity_source=console_authz.IDENTITY_ACCESS_PROXY,
authenticated=True,
)
def viewer() -> console_authz.Principal:
return console_authz.Principal(
subject="[email protected]",
role=console_authz.VIEWER,
identity_source=console_authz.IDENTITY_ACCESS_PROXY,
authenticated=True,
)
class TestCapabilityWiring(unittest.TestCase):
"""AC1: authority is declared, not invented by the console."""
def test_actions_resolve_through_the_capability_map(self):
for action_id in (
sanctioned_restart.ACTION_RESTART_NAMESPACE,
sanctioned_restart.ACTION_RELOAD_NAMESPACE,
):
with self.subTest(action=action_id):
action = console_authz.get_action(action_id)
self.assertIsNotNone(action)
self.assertIn(action.task_key, TASK_CAPABILITY_MAP)
self.assertEqual(
action.mcp_permission,
TASK_CAPABILITY_MAP[action.task_key]["permission"],
)
def test_restart_permission_is_not_a_gitea_operation(self):
"""No configured Gitea profile should satisfy a host restart."""
permission = TASK_CAPABILITY_MAP["restart_namespace"]["permission"]
self.assertFalse(permission.startswith("gitea."))
def test_restart_is_destructive_dual_control_break_glass(self):
action = console_authz.get_action(
sanctioned_restart.ACTION_RESTART_NAMESPACE
)
self.assertEqual(action.action_class, console_authz.CLASS_DESTRUCTIVE)
self.assertEqual(action.minimum_role, console_authz.ADMIN)
self.assertTrue(action.dual_control)
self.assertTrue(action.break_glass)
self.assertTrue(action.requires_confirmation)
def test_reload_is_privileged_but_not_destructive(self):
action = console_authz.get_action(
sanctioned_restart.ACTION_RELOAD_NAMESPACE
)
self.assertEqual(action.action_class, console_authz.CLASS_PRIVILEGED)
self.assertTrue(action.requires_confirmation)
class TestPreview(unittest.TestCase):
"""AC4: an authorized preview renders the plan without executing it."""
def test_preview_lists_the_mutation_ledger(self):
preview = sanctioned_restart.build_restart_preview(
NAMESPACE, principal=admin(), env=READY_ENV
)
steps = [entry["step"] for entry in preview["mutation_ledger"]]
self.assertEqual(
steps, ["quiesce", "host_restart_hook", "health_recheck", "audit"]
)
self.assertTrue(preview["scope_valid"])
self.assertTrue(preview["post_restart_verification_required"])
def test_reload_preview_drains_instead_of_restarting(self):
preview = sanctioned_restart.build_restart_preview(
NAMESPACE, sanctioned_restart.MODE_RELOAD,
principal=admin(), env=READY_ENV,
)
steps = [entry["step"] for entry in preview["mutation_ledger"]]
self.assertIn("host_graceful_reload", steps)
self.assertNotIn("host_restart_hook", steps)
def test_preview_never_enables_execution(self):
preview = sanctioned_restart.build_restart_preview(
NAMESPACE, principal=admin(), env=READY_ENV
)
self.assertFalse(preview["execution_enabled"])
self.assertFalse(preview["authorization"]["execution_enabled"])
def test_confirmation_phrase_binds_the_namespace(self):
self.assertTrue(
sanctioned_restart.confirmation_matches(
NAMESPACE, sanctioned_restart.MODE_RESTART,
"restart gitea-author",
)
)
# A phrase typed for one namespace must not authorize another.
self.assertFalse(
sanctioned_restart.confirmation_matches(
"gitea-merger", sanctioned_restart.MODE_RESTART,
"restart gitea-author",
)
)
class TestGates(unittest.TestCase):
"""AC1/AC4: every gate denies with a stable reason code."""
def _assess(self, **kwargs):
params = {
"principal": admin(),
"confirmation": f"restart {NAMESPACE}",
"env": READY_ENV,
}
params.update(kwargs)
namespace = params.pop("namespace", NAMESPACE)
mode = params.pop("mode", sanctioned_restart.MODE_RESTART)
return sanctioned_restart.assess_restart_request(
namespace, mode, **params
)
def test_authorized_confirmed_request_passes_every_gate(self):
result = self._assess()
self.assertTrue(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.ALLOW_HOST_ACTION_REQUIRED
)
def test_passing_every_gate_is_not_an_execution_grant(self):
"""An allowed request still never lets the console touch the process."""
result = self._assess()
self.assertTrue(result["allowed"])
self.assertFalse(result["execution_enabled"])
self.assertFalse(result["console_executes"])
def test_unauthorized_principal_is_denied(self):
result = self._assess(principal=viewer())
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_UNAUTHORIZED
)
def test_anonymous_principal_is_denied(self):
result = self._assess(principal=None)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_UNAUTHORIZED
)
def test_missing_confirmation_is_denied(self):
result = self._assess(confirmation=None)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_CONFIRMATION_MISSING
)
def test_confirmation_for_another_namespace_is_denied(self):
result = self._assess(confirmation="restart gitea-merger")
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_CONFIRMATION_MISMATCH
)
def test_missing_operator_authorization_is_denied(self):
env = {sanctioned_restart.RESTART_HOOK_ENV: "launchd:cc.prgs.author"}
result = self._assess(env=env)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"],
sanctioned_restart.DENY_OPERATOR_AUTHORIZATION,
)
def test_missing_host_hook_is_denied_without_kill_fallback(self):
env = {
runtime_recovery_guard.OPERATOR_AUTHORIZATION_ENV: "ops-ticket-1",
}
result = self._assess(env=env)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_HOOK_NOT_CONFIGURED
)
def test_fleet_scope_is_refused(self):
for scope in ("all", "*", "fleet"):
with self.subTest(scope=scope):
result = self._assess(
namespace=scope, confirmation=f"restart {scope}"
)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_FLEET_SCOPE
)
def test_unknown_namespace_is_refused(self):
result = self._assess(
namespace="gitea-nope", confirmation="restart gitea-nope"
)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_UNKNOWN_NAMESPACE
)
def test_unknown_mode_is_refused(self):
result = self._assess(mode="obliterate")
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_UNKNOWN_MODE
)
def test_live_contamination_marker_blocks_restart(self):
marker = runtime_recovery_guard.build_contamination_record(
reason_class=runtime_recovery_guard.REASON_MANUAL_DAEMON_KILL,
command_redacted="pkill -f mcp_server.py",
)
result = self._assess(contamination_marker=marker)
self.assertFalse(result["allowed"])
self.assertEqual(
result["reason_code"], sanctioned_restart.DENY_CONTAMINATED_RUNTIME
)
def test_reconciler_cleared_marker_no_longer_blocks(self):
marker = runtime_recovery_guard.build_contamination_record(
reason_class=runtime_recovery_guard.REASON_MANUAL_DAEMON_KILL,
command_redacted="pkill -f mcp_server.py",
)
marker = dict(marker, cleared_by_reconciler=True)
result = self._assess(contamination_marker=marker)
self.assertTrue(result["allowed"])
class TestExecutionNeverKills(unittest.TestCase):
"""AC5: no path exposes or runs a raw process kill."""
def test_authorized_execution_defers_to_the_host_supervisor(self):
result = sanctioned_restart.execute_restart(
NAMESPACE,
principal=admin(),
confirmation=f"restart {NAMESPACE}",
env=READY_ENV,
)
self.assertTrue(result["allowed"])
self.assertFalse(result["success"])
self.assertFalse(result["process_kill_executed"])
self.assertEqual(
result["outcome"], sanctioned_restart.ALLOW_HOST_ACTION_REQUIRED
)
def test_denied_execution_reports_the_refusing_gate(self):
result = sanctioned_restart.execute_restart(
NAMESPACE, principal=viewer(), confirmation=f"restart {NAMESPACE}",
env=READY_ENV,
)
self.assertFalse(result["allowed"])
self.assertEqual(
result["outcome"], sanctioned_restart.DENY_UNAUTHORIZED
)
self.assertFalse(result["process_kill_executed"])
def test_module_never_spawns_a_process(self):
path = os.path.join(
os.path.dirname(os.path.dirname(os.path.abspath(__file__))),
"webui", "sanctioned_restart.py",
)
with open(path, encoding="utf-8") as handle:
source = handle.read()
for forbidden in (
"import subprocess", "import signal", "os.kill", "os.system",
"popen",
):
with self.subTest(forbidden=forbidden):
self.assertNotIn(forbidden, source.lower())
def test_no_surface_returns_a_kill_command(self):
payloads = [
sanctioned_restart.build_restart_preview(
NAMESPACE, principal=admin(), env=READY_ENV
),
sanctioned_restart.restart_policy(),
sanctioned_restart.execute_restart(
NAMESPACE, principal=admin(),
confirmation=f"restart {NAMESPACE}", env=READY_ENV,
),
]
for payload in payloads:
rendered = json.dumps(payload, default=str).lower()
self.assertNotIn("kill -9", rendered)
self.assertNotIn("pkill -f", rendered)
def test_policy_declares_no_raw_kill_and_no_silent_restart(self):
policy = sanctioned_restart.restart_policy()
self.assertFalse(policy["raw_kill_exposed"])
self.assertFalse(policy["console_executes_process_kill"])
self.assertFalse(policy["fleet_scope_permitted"])
self.assertFalse(policy["silent_auto_restart_permitted"])
self.assertTrue(policy["audit_required"])
class TestContaminationClassification(unittest.TestCase):
"""AC2: manual pkill is contamination, and it blocks clean claims."""
def test_manual_daemon_pkill_is_contamination(self):
result = sanctioned_restart.classify_restart_command(
"pkill -f mcp_server.py"
)
self.assertTrue(result["contamination"])
self.assertFalse(result["clean_claim_allowed"])
self.assertIsNotNone(result["contamination_marker"])
self.assertEqual(
result["sanctioned_alternative"],
sanctioned_restart.ACTION_RESTART_NAMESPACE,
)
def test_broad_process_kill_is_contamination(self):
result = sanctioned_restart.classify_restart_command("killall -9 Python")
self.assertTrue(result["contamination"])
self.assertFalse(result["clean_claim_allowed"])
def test_marker_names_the_sanctioned_alternative(self):
result = sanctioned_restart.classify_restart_command(
"pkill -f mcp_server.py"
)
marker = result["contamination_marker"]
self.assertIn(
sanctioned_restart.ACTION_RESTART_NAMESPACE, marker["detail"]
)
def test_benign_command_is_not_contamination(self):
result = sanctioned_restart.classify_restart_command("git status")
self.assertFalse(result["contamination"])
self.assertTrue(result["clean_claim_allowed"])
def test_no_command_is_not_contamination(self):
result = sanctioned_restart.classify_restart_command(None)
self.assertFalse(result["contamination"])
self.assertTrue(result["clean_claim_allowed"])
class TestPostRestartHealth(unittest.TestCase):
"""AC3: a clean post-restart claim needs live client-namespace proof."""
def test_live_client_probe_clears_the_session(self):
result = sanctioned_restart.verify_post_restart_health(
NAMESPACE,
probe_result={"success": True},
probe_source=mcp_namespace_health.PROBE_SOURCE_CLIENT,
registered_tools=["gitea_whoami"],
required_tool="gitea_whoami",
)
self.assertEqual(result["status"], sanctioned_restart.HEALTH_CLEAN)
self.assertTrue(result["clean_claim_allowed"])
self.assertTrue(result["mutations_allowed"])
def test_offline_probe_does_not_clear_the_session(self):
result = sanctioned_restart.verify_post_restart_health(
NAMESPACE,
probe_result={"success": True},
probe_source=mcp_namespace_health.PROBE_SOURCE_OFFLINE,
registered_tools=["gitea_whoami"],
required_tool="gitea_whoami",
)
self.assertFalse(result["clean_claim_allowed"])
self.assertFalse(result["mutations_allowed"])
def test_failed_probe_is_unhealthy(self):
result = sanctioned_restart.verify_post_restart_health(
NAMESPACE,
probe_result={"success": False, "error": "client is closing: EOF"},
probe_source=mcp_namespace_health.PROBE_SOURCE_CLIENT,
registered_tools=["gitea_whoami"],
required_tool="gitea_whoami",
)
self.assertEqual(result["status"], sanctioned_restart.HEALTH_UNHEALTHY)
self.assertFalse(result["clean_claim_allowed"])
def test_static_registration_alone_never_clears_the_session(self):
result = sanctioned_restart.verify_post_restart_health(
NAMESPACE,
registered_tools=["gitea_whoami"],
required_tool="gitea_whoami",
)
self.assertFalse(result["clean_claim_allowed"])
class TestAuditEmission(unittest.TestCase):
"""Every restart attempt is audited with actor, target, and result."""
def _run(self, principal, sink):
prior = os.environ.get(console_audit.AUDIT_LOG_ENV)
os.environ[console_audit.AUDIT_LOG_ENV] = sink
try:
return sanctioned_restart.execute_restart(
NAMESPACE,
principal=principal,
confirmation=f"restart {NAMESPACE}",
env=READY_ENV,
request_id="req-642",
)
finally:
if prior is None:
os.environ.pop(console_audit.AUDIT_LOG_ENV, None)
else:
os.environ[console_audit.AUDIT_LOG_ENV] = prior
def test_allowed_attempt_is_written_with_actor_and_target(self):
with tempfile.TemporaryDirectory() as tmp:
sink = os.path.join(tmp, "audit.jsonl")
result = self._run(admin(), sink)
self.assertTrue(result["audit"]["written"])
with open(sink, encoding="utf-8") as handle:
record = json.loads(handle.read().strip())
self.assertEqual(
record["action"], sanctioned_restart.ACTION_RESTART_NAMESPACE
)
self.assertEqual(record["target"]["namespace"], NAMESPACE)
self.assertEqual(record["target"]["mode"], "restart")
self.assertEqual(record["result"], console_audit.RESULT_ALLOWED)
self.assertEqual(record["actor"]["subject"], "[email protected]")
self.assertFalse(record["metadata"]["process_kill_executed"])
def test_denied_attempt_is_audited_too(self):
with tempfile.TemporaryDirectory() as tmp:
sink = os.path.join(tmp, "audit.jsonl")
self._run(viewer(), sink)
with open(sink, encoding="utf-8") as handle:
record = json.loads(handle.read().strip())
self.assertEqual(record["result"], console_audit.RESULT_DENIED)
self.assertEqual(
record["reason_code"], sanctioned_restart.DENY_UNAUTHORIZED
)
def test_restart_audit_uses_break_glass_retention(self):
action = console_authz.get_action(
sanctioned_restart.ACTION_RESTART_NAMESPACE
)
self.assertEqual(
console_audit.retention_class_for(action),
console_audit.RETENTION_BREAK_GLASS,
)
class TestRegistrySurface(unittest.TestCase):
"""AC5: the console surfaces the control, still disabled, with no kill."""
def test_registry_exposes_both_actions_disabled(self):
registry = gated_actions.load_action_registry()
for action_id in (
sanctioned_restart.ACTION_RESTART_NAMESPACE,
sanctioned_restart.ACTION_RELOAD_NAMESPACE,
):
with self.subTest(action=action_id):
action = registry.get(action_id)
self.assertIsNotNone(action)
self.assertFalse(action.enabled)
def test_registry_preview_names_the_namespace_target(self):
preview = gated_actions.preview_action(
sanctioned_restart.ACTION_RESTART_NAMESPACE, namespace=NAMESPACE
)
target = preview["mutation_ledger"][0]["target"]
self.assertIn(NAMESPACE, target)
self.assertFalse(preview["enabled"])
def test_registry_attempt_fails_closed(self):
result = gated_actions.attempt_action(
sanctioned_restart.ACTION_RESTART_NAMESPACE, namespace=NAMESPACE
)
self.assertFalse(result["success"])
if __name__ == "__main__":
unittest.main()
+433
View File
@@ -0,0 +1,433 @@
"""Model usage, token cost, latency, and workflow-performance analytics (#651, Phase 4).
Ingests session instrumentation metrics, aggregates usage/cost/latency percentiles
by project, role, model, issue/PR, and stage, enforcing secret redaction and
explicitly rendering missing metrics as "Unknown" without zero-fabrication.
"""
from __future__ import annotations
import math
from dataclasses import asdict, dataclass
from typing import Any, Sequence
import control_plane_db
from webui import console_redaction
ANALYTICS_SCHEMA_VERSION = 1
@dataclass(frozen=True)
class UsageEvent:
usage_id: int
session_id: str | None
remote: str
org: str
repo: str
project_id: str | None
role: str
model: str
issue_number: int | None
pr_number: int | None
stage: str
input_tokens: int | None
output_tokens: int | None
total_tokens: int | None
estimated_cost_usd: float | None
latency_ms: int | None
duration_ms: int | None
status: str
metadata: str | None
created_at: str
def to_dict(self) -> dict[str, Any]:
d = asdict(self)
if d["metadata"]:
d["metadata"] = console_redaction.redact_text(d["metadata"])
return d
@dataclass(frozen=True)
class GroupMetrics:
name: str
total_events: int
events_with_tokens: int
input_tokens: int | None
output_tokens: int | None
total_tokens: int | None
events_with_cost: int
estimated_cost_usd: float | None
events_with_latency: int
latency_p50_ms: float | None
latency_p90_ms: float | None
latency_p95_ms: float | None
latency_p99_ms: float | None
latency_avg_ms: float | None
events_with_duration: int
duration_avg_ms: float | None
display_tokens: str
display_cost: str
display_latency_p50: str
display_latency_p90: str
display_duration_avg: str
def to_dict(self) -> dict[str, Any]:
return asdict(self)
@dataclass(frozen=True)
class AnalyticsSnapshot:
ok: bool
reason: str
schema_version: int
remote: str
org: str
repo: str
total_events: int
overall_summary: GroupMetrics
by_project: dict[str, GroupMetrics]
by_role: dict[str, GroupMetrics]
by_model: dict[str, GroupMetrics]
by_work_item: dict[str, GroupMetrics]
by_stage: dict[str, GroupMetrics]
events: tuple[UsageEvent, ...]
def to_dict(self) -> dict[str, Any]:
return {
"ok": self.ok,
"reason": self.reason,
"schema_version": self.schema_version,
"remote": self.remote,
"org": self.org,
"repo": self.repo,
"total_events": self.total_events,
"overall_summary": self.overall_summary.to_dict(),
"by_project": {k: v.to_dict() for k, v in self.by_project.items()},
"by_role": {k: v.to_dict() for k, v in self.by_role.items()},
"by_model": {k: v.to_dict() for k, v in self.by_model.items()},
"by_work_item": {k: v.to_dict() for k, v in self.by_work_item.items()},
"by_stage": {k: v.to_dict() for k, v in self.by_stage.items()},
"events": [e.to_dict() for e in self.events],
}
def compute_percentile(values: Sequence[float | int], percentile: float) -> float | None:
if not values:
return None
sorted_vals = sorted(values)
n = len(sorted_vals)
if n == 1:
return float(sorted_vals[0])
k = (n - 1) * (percentile / 100.0)
f = math.floor(k)
c = math.ceil(k)
if f == c:
return float(sorted_vals[int(f)])
d0 = sorted_vals[int(f)] * (c - k)
d1 = sorted_vals[int(c)] * (k - f)
return float(d0 + d1)
def aggregate_events(group_name: str, events: Sequence[UsageEvent]) -> GroupMetrics:
total_events = len(events)
if total_events == 0:
return GroupMetrics(
name=group_name,
total_events=0,
events_with_tokens=0,
input_tokens=None,
output_tokens=None,
total_tokens=None,
events_with_cost=0,
estimated_cost_usd=None,
events_with_latency=0,
latency_p50_ms=None,
latency_p90_ms=None,
latency_p95_ms=None,
latency_p99_ms=None,
latency_avg_ms=None,
events_with_duration=0,
duration_avg_ms=None,
display_tokens="Unknown",
display_cost="Unknown",
display_latency_p50="Unknown",
display_latency_p90="Unknown",
display_duration_avg="Unknown",
)
token_events = [
e for e in events
if e.total_tokens is not None or e.input_tokens is not None or e.output_tokens is not None
]
events_with_tokens = len(token_events)
if events_with_tokens > 0:
input_tokens = sum(e.input_tokens or 0 for e in token_events)
output_tokens = sum(e.output_tokens or 0 for e in token_events)
total_tokens = sum(
e.total_tokens if e.total_tokens is not None else ((e.input_tokens or 0) + (e.output_tokens or 0))
for e in token_events
)
display_tokens = f"{total_tokens:,}"
else:
input_tokens = None
output_tokens = None
total_tokens = None
display_tokens = "Unknown"
cost_events = [e for e in events if e.estimated_cost_usd is not None]
events_with_cost = len(cost_events)
if events_with_cost > 0:
estimated_cost_usd = round(sum(e.estimated_cost_usd for e in cost_events), 6)
display_cost = f"${estimated_cost_usd:.4f}"
else:
estimated_cost_usd = None
display_cost = "Unknown"
latency_vals = [e.latency_ms for e in events if e.latency_ms is not None]
events_with_latency = len(latency_vals)
if events_with_latency > 0:
latency_p50_ms = compute_percentile(latency_vals, 50.0)
latency_p90_ms = compute_percentile(latency_vals, 90.0)
latency_p95_ms = compute_percentile(latency_vals, 95.0)
latency_p99_ms = compute_percentile(latency_vals, 99.0)
latency_avg_ms = round(sum(latency_vals) / events_with_latency, 2)
display_latency_p50 = f"{round(latency_p50_ms, 1)} ms" if latency_p50_ms is not None else "Unknown"
display_latency_p90 = f"{round(latency_p90_ms, 1)} ms" if latency_p90_ms is not None else "Unknown"
else:
latency_p50_ms = None
latency_p90_ms = None
latency_p95_ms = None
latency_p99_ms = None
latency_avg_ms = None
display_latency_p50 = "Unknown"
display_latency_p90 = "Unknown"
duration_vals = [e.duration_ms for e in events if e.duration_ms is not None]
events_with_duration = len(duration_vals)
if events_with_duration > 0:
duration_avg_ms = round(sum(duration_vals) / events_with_duration, 2)
display_duration_avg = f"{round(duration_avg_ms / 1000.0, 2)} s" if duration_avg_ms >= 1000 else f"{round(duration_avg_ms, 1)} ms"
else:
duration_avg_ms = None
display_duration_avg = "Unknown"
return GroupMetrics(
name=group_name,
total_events=total_events,
events_with_tokens=events_with_tokens,
input_tokens=input_tokens,
output_tokens=output_tokens,
total_tokens=total_tokens,
events_with_cost=events_with_cost,
estimated_cost_usd=estimated_cost_usd,
events_with_latency=events_with_latency,
latency_p50_ms=latency_p50_ms,
latency_p90_ms=latency_p90_ms,
latency_p95_ms=latency_p95_ms,
latency_p99_ms=latency_p99_ms,
latency_avg_ms=latency_avg_ms,
events_with_duration=events_with_duration,
duration_avg_ms=duration_avg_ms,
display_tokens=display_tokens,
display_cost=display_cost,
display_latency_p50=display_latency_p50,
display_latency_p90=display_latency_p90,
display_duration_avg=display_duration_avg,
)
def record_usage(
*,
db_path: str | None = None,
session_id: str | None = None,
remote: str = "dadeschools",
org: str = "",
repo: str = "",
project_id: str | None = None,
role: str = "unknown",
model: str = "unknown",
issue_number: int | None = None,
pr_number: int | None = None,
stage: str = "unknown",
input_tokens: int | None = None,
output_tokens: int | None = None,
total_tokens: int | None = None,
estimated_cost_usd: float | None = None,
latency_ms: int | None = None,
duration_ms: int | None = None,
status: str = "success",
metadata: str | dict[str, Any] | None = None,
created_at: str | None = None,
) -> int:
"""Ingest/record a single usage event with optional metrics."""
db = control_plane_db.ControlPlaneDB(db_path=db_path)
return db.record_usage_event(
session_id=session_id,
remote=remote,
org=org,
repo=repo,
project_id=project_id,
role=role,
model=model,
issue_number=issue_number,
pr_number=pr_number,
stage=stage,
input_tokens=input_tokens,
output_tokens=output_tokens,
total_tokens=total_tokens,
estimated_cost_usd=estimated_cost_usd,
latency_ms=latency_ms,
duration_ms=duration_ms,
status=status,
metadata=metadata,
created_at=created_at,
)
def load_analytics(
*,
db_path: str | None = None,
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
project_id: str | None = None,
role: str | None = None,
model: str | None = None,
stage: str | None = None,
issue_number: int | None = None,
pr_number: int | None = None,
limit: int = 500,
) -> AnalyticsSnapshot:
"""Load analytics snapshot aggregated by project, role, model, issue/PR, and stage."""
remote_filter = (remote or "").strip() or None
org_filter = (org or "").strip() or None
repo_filter = (repo or "").strip() or None
role_filter = (role or "").strip() or None
model_filter = (model or "").strip() or None
stage_filter = (stage or "").strip() or None
# F4: coerce optional scope filters to str so AnalyticsSnapshot never holds None.
scope_remote = (remote or "").strip()
scope_org = (org or "").strip()
scope_repo = (repo or "").strip()
try:
db = control_plane_db.ControlPlaneDB(db_path=db_path)
rows = db.query_usage_events(
remote=remote_filter,
org=org_filter,
repo=repo_filter,
project_id=project_id,
role=role_filter,
model=model_filter,
stage=stage_filter,
issue_number=issue_number,
pr_number=pr_number,
limit=limit,
)
except Exception as exc:
empty_summary = aggregate_events("Overall", [])
return AnalyticsSnapshot(
ok=False,
reason=f"control_plane_db_unavailable: {exc}",
schema_version=ANALYTICS_SCHEMA_VERSION,
remote=scope_remote,
org=scope_org,
repo=scope_repo,
total_events=0,
overall_summary=empty_summary,
by_project={},
by_role={},
by_model={},
by_work_item={},
by_stage={},
events=(),
)
parsed_events: list[UsageEvent] = []
for r in rows:
meta = console_redaction.redact_text(r.get("metadata")) if r.get("metadata") else None
parsed_events.append(
UsageEvent(
usage_id=r["usage_id"],
session_id=r.get("session_id"),
remote=r.get("remote") or scope_remote,
org=r.get("org") or scope_org,
repo=r.get("repo") or scope_repo,
project_id=r.get("project_id"),
role=r.get("role") or "unknown",
model=r.get("model") or "unknown",
issue_number=r.get("issue_number"),
pr_number=r.get("pr_number"),
stage=r.get("stage") or "unknown",
input_tokens=r.get("input_tokens"),
output_tokens=r.get("output_tokens"),
total_tokens=r.get("total_tokens"),
estimated_cost_usd=r.get("estimated_cost_usd"),
latency_ms=r.get("latency_ms"),
duration_ms=r.get("duration_ms"),
status=r.get("status") or "success",
metadata=meta,
created_at=r.get("created_at") or "",
)
)
overall_summary = aggregate_events("Overall", parsed_events)
# Group by project
groups_by_project: dict[str, list[UsageEvent]] = {}
for e in parsed_events:
key = e.project_id or (f"{e.org}/{e.repo}" if e.org and e.repo else "default")
groups_by_project.setdefault(key, []).append(e)
by_project = {k: aggregate_events(k, v) for k, v in groups_by_project.items()}
# Group by role
groups_by_role: dict[str, list[UsageEvent]] = {}
for e in parsed_events:
groups_by_role.setdefault(e.role, []).append(e)
by_role = {k: aggregate_events(k, v) for k, v in groups_by_role.items()}
# Group by model
groups_by_model: dict[str, list[UsageEvent]] = {}
for e in parsed_events:
groups_by_model.setdefault(e.model, []).append(e)
by_model = {k: aggregate_events(k, v) for k, v in groups_by_model.items()}
# Group by work item
groups_by_work_item: dict[str, list[UsageEvent]] = {}
for e in parsed_events:
if e.issue_number:
key = f"issue #{e.issue_number}"
elif e.pr_number:
key = f"pr #{e.pr_number}"
else:
key = "unlinked"
groups_by_work_item.setdefault(key, []).append(e)
by_work_item = {k: aggregate_events(k, v) for k, v in groups_by_work_item.items()}
# Group by stage
groups_by_stage: dict[str, list[UsageEvent]] = {}
for e in parsed_events:
groups_by_stage.setdefault(e.stage, []).append(e)
by_stage = {k: aggregate_events(k, v) for k, v in groups_by_stage.items()}
return AnalyticsSnapshot(
ok=True,
reason="ok",
schema_version=ANALYTICS_SCHEMA_VERSION,
remote=scope_remote,
org=scope_org,
repo=scope_repo,
total_events=len(parsed_events),
overall_summary=overall_summary,
by_project=by_project,
by_role=by_role,
by_model=by_model,
by_work_item=by_work_item,
by_stage=by_stage,
events=tuple(parsed_events),
)
def snapshot_to_dict(snapshot: AnalyticsSnapshot) -> dict[str, Any]:
return snapshot.to_dict()
+248
View File
@@ -0,0 +1,248 @@
"""HTML views for the Model Usage & Performance Analytics console (#651)."""
from __future__ import annotations
import html
from webui.analytics_loader import AnalyticsSnapshot, GroupMetrics, UsageEvent
from webui.layout import render_page
def _escape(text: object) -> str:
"""HTML-escape dynamic analytics fields (mirrors audit_views / project_views)."""
return html.escape(str(text), quote=True)
def _render_badge(text: str, badge_type: str = "muted") -> str:
return f'<span class="badge badge-{_escape(badge_type)}">{_escape(text)}</span>'
def _render_group_table(title: str, groups: dict[str, GroupMetrics], key_header: str = "Group") -> str:
if not groups:
return (
f"<h3>{_escape(title)}</h3>"
'<div class="card"><p class="muted">No telemetry events recorded for this dimension.</p></div>'
)
rows = []
for key, g in sorted(groups.items(), key=lambda x: x[1].total_events, reverse=True):
cost_cell = (
f'<span class="accent">{_escape(g.display_cost)}</span>'
if g.events_with_cost > 0
else _render_badge("Unknown")
)
tokens_cell = (
_escape(g.display_tokens)
if g.events_with_tokens > 0
else _render_badge("Unknown")
)
lat_p50 = (
_escape(g.display_latency_p50)
if g.events_with_latency > 0
else _render_badge("Unknown")
)
lat_p90 = (
_escape(g.display_latency_p90)
if g.events_with_latency > 0
else _render_badge("Unknown")
)
dur_avg = (
_escape(g.display_duration_avg)
if g.events_with_duration > 0
else _render_badge("Unknown")
)
rows.append(
"<tr>"
f"<td><strong>{_escape(key)}</strong></td>"
f"<td>{g.total_events}</td>"
f"<td>{tokens_cell}</td>"
f"<td>{cost_cell}</td>"
f"<td>{lat_p50}</td>"
f"<td>{lat_p90}</td>"
f"<td>{dur_avg}</td>"
"</tr>"
)
rows_html = "".join(rows)
return f"""
<h3>{_escape(title)}</h3>
<div class="card" style="overflow-x: auto;">
<table class="data-table">
<thead>
<tr>
<th>{_escape(key_header)}</th>
<th>Events</th>
<th>Total Tokens</th>
<th>Est. Cost</th>
<th>Latency (p50)</th>
<th>Latency (p90)</th>
<th>Avg Stage Duration</th>
</tr>
</thead>
<tbody>
{rows_html}
</tbody>
</table>
</div>
"""
def _render_events_table(events: tuple[UsageEvent, ...]) -> str:
if not events:
return (
"<h3>Recent Usage & Instrumentation Events</h3>"
'<div class="card"><p class="muted">No individual telemetry events recorded yet. Opt-in instrumentation via session logging or authorized POST /api/v1/analytics/usage.</p></div>'
)
rows = []
for e in list(events)[-50:]: # Display latest 50
if e.issue_number is not None:
work_item = f"issue #{e.issue_number}"
elif e.pr_number is not None:
work_item = f"pr #{e.pr_number}"
else:
work_item = "unlinked"
tokens = (
_escape(f"{e.total_tokens:,}")
if e.total_tokens is not None
else _render_badge("Unknown")
)
cost = (
_escape(f"${e.estimated_cost_usd:.4f}")
if e.estimated_cost_usd is not None
else _render_badge("Unknown")
)
latency = (
_escape(f"{e.latency_ms} ms")
if e.latency_ms is not None
else _render_badge("Unknown")
)
duration = (
_escape(f"{e.duration_ms} ms")
if e.duration_ms is not None
else _render_badge("Unknown")
)
status_badge = _render_badge(
e.status, "success" if e.status == "success" else "danger"
)
rows.append(
"<tr>"
f"<td>#{e.usage_id}</td>"
f"<td><small>{_escape(e.created_at)}</small></td>"
f"<td><span class=\"badge\">{_escape(e.role)}</span></td>"
f"<td><strong>{_escape(e.model)}</strong></td>"
f"<td>{_escape(e.stage)}</td>"
f"<td>{_escape(work_item)}</td>"
f"<td>{tokens}</td>"
f"<td>{cost}</td>"
f"<td>{latency}</td>"
f"<td>{duration}</td>"
f"<td>{status_badge}</td>"
"</tr>"
)
rows_html = "".join(rows)
return f"""
<h3>Recent Telemetry Events</h3>
<div class="card" style="overflow-x: auto;">
<table class="data-table">
<thead>
<tr>
<th>ID</th>
<th>Timestamp</th>
<th>Role</th>
<th>Model</th>
<th>Stage</th>
<th>Work Item</th>
<th>Tokens</th>
<th>Cost</th>
<th>Latency</th>
<th>Duration</th>
<th>Status</th>
</tr>
</thead>
<tbody>
{rows_html}
</tbody>
</table>
</div>
"""
def render_analytics_page(snapshot: AnalyticsSnapshot) -> str:
"""Render the main Model Usage & Performance Analytics console page."""
summary = snapshot.overall_summary
kpi_tokens = (
_escape(summary.display_tokens)
if summary.events_with_tokens > 0
else _render_badge("Unknown")
)
kpi_cost = (
_escape(summary.display_cost)
if summary.events_with_cost > 0
else _render_badge("Unknown")
)
kpi_lat_p50 = (
_escape(summary.display_latency_p50)
if summary.events_with_latency > 0
else _render_badge("Unknown")
)
kpi_dur_avg = (
_escape(summary.display_duration_avg)
if summary.events_with_duration > 0
else _render_badge("Unknown")
)
status_notice = ""
if not snapshot.ok:
status_notice = (
f'<div class="card warning-card"><strong>Degraded Data Source:</strong> '
f'{_escape(snapshot.reason)}</div>'
)
body_html = f"""
<h2>Model Usage & Performance Analytics (Phase 4)</h2>
<p class="muted">
Durable console analytics for model usage, token cost, latency percentiles, and workflow-stage performance correlated to issues, PRs, and worker roles.
</p>
{status_notice}
<div class="notice-card" style="background: rgba(91, 159, 212, 0.1); border: 1px solid var(--border); padding: 0.75rem 1rem; border-radius: 6px; margin-bottom: 1.5rem;">
<small><strong>Note on telemetry fidelity:</strong> Missing data or untracked metrics are explicitly labeled as <em>Unknown</em>. No token costs or latency metrics are zero-fabricated.</small>
</div>
<div class="card-grid" style="display: grid; grid-template-columns: repeat(auto-fit, minmax(180px, 1fr)); gap: 1rem; margin-bottom: 1.5rem;">
<div class="card">
<span class="muted" style="font-size: 0.85rem;">Total Events</span>
<h3 style="margin: 0.25rem 0 0 0;">{summary.total_events}</h3>
</div>
<div class="card">
<span class="muted" style="font-size: 0.85rem;">Total Tokens</span>
<h3 style="margin: 0.25rem 0 0 0;">{kpi_tokens}</h3>
</div>
<div class="card">
<span class="muted" style="font-size: 0.85rem;">Est. Token Cost</span>
<h3 style="margin: 0.25rem 0 0 0;">{kpi_cost}</h3>
</div>
<div class="card">
<span class="muted" style="font-size: 0.85rem;">Latency (p50)</span>
<h3 style="margin: 0.25rem 0 0 0;">{kpi_lat_p50}</h3>
</div>
<div class="card">
<span class="muted" style="font-size: 0.85rem;">Avg Stage Duration</span>
<h3 style="margin: 0.25rem 0 0 0;">{kpi_dur_avg}</h3>
</div>
</div>
{_render_group_table("Usage & Cost by Model", snapshot.by_model, "Model")}
{_render_group_table("Performance by Workflow Stage", snapshot.by_stage, "Stage")}
{_render_group_table("Usage & Cost by Role", snapshot.by_role, "Role")}
{_render_group_table("Work Item Analytics", snapshot.by_work_item, "Work Item")}
{_render_events_table(snapshot.events)}
"""
return render_page(title="Model Usage & Performance Analytics", body_html=body_html)
+118
View File
@@ -47,6 +47,12 @@ from webui.worktree_views import render_worktrees_page
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
from webui.runtime_views import render_runtime_page
from webui.timeline import load_timeline, snapshot_to_dict as timeline_snapshot_to_dict
from webui.analytics_loader import (
load_analytics,
record_usage,
snapshot_to_dict as analytics_snapshot_to_dict,
)
from webui.analytics_views import render_analytics_page
from webui.system_health import (
API_PATH as SYSTEM_HEALTH_API_PATH,
load_system_health,
@@ -567,6 +573,114 @@ async def api_v1_timeline(request: Request) -> JSONResponse:
return JSONResponse(timeline_snapshot_to_dict(snapshot), status_code=status_code)
async def analytics(request: Request) -> HTMLResponse:
"""Read-only model usage, token cost, latency, and performance analytics HTML view (#651)."""
snapshot = load_analytics(
remote=request.query_params.get("remote"),
org=request.query_params.get("org"),
repo=request.query_params.get("repo"),
role=request.query_params.get("role"),
model=request.query_params.get("model"),
stage=request.query_params.get("stage"),
issue_number=_query_int(request, "issue"),
pr_number=_query_int(request, "pr"),
limit=_query_int(request, "limit") or 200,
)
return HTMLResponse(render_analytics_page(snapshot))
async def api_v1_analytics(request: Request) -> JSONResponse:
"""Read-only model usage, token cost, latency, and performance analytics API (#651)."""
snapshot = load_analytics(
remote=request.query_params.get("remote"),
org=request.query_params.get("org"),
repo=request.query_params.get("repo"),
role=request.query_params.get("role"),
model=request.query_params.get("model"),
stage=request.query_params.get("stage"),
issue_number=_query_int(request, "issue"),
pr_number=_query_int(request, "pr"),
limit=_query_int(request, "limit") or 500,
)
status_code = 200 if snapshot.ok else 500
return JSONResponse(analytics_snapshot_to_dict(snapshot), status_code=status_code)
async def api_v1_analytics_ingest(request: Request) -> JSONResponse:
"""Optional session instrumentation ingestion endpoint (#651).
Fail-closed write: every request is authorized through console_authz
(``record_analytics_usage``) before any control-plane DB mutation. Phase 1
keeps ``execution_enabled=False`` and denies unauthenticated callers, so
this route cannot be used as an unauthenticated write or XSS injection
vector (PR #876 F2).
"""
try:
body = await request.json()
except Exception:
body = {}
if not isinstance(body, dict):
body = {}
principal = resolve_principal(headers=dict(request.headers))
decision = authorize(
"record_analytics_usage", principal, for_execution=True
)
allowed = bool(decision.allowed and decision.execution_enabled)
console_audit.record_event(
action_id="record_analytics_usage",
result=(
console_audit.RESULT_ALLOWED
if allowed
else console_audit.RESULT_DENIED
),
decision=decision,
principal=principal,
target=_audit_target("record_analytics_usage", body),
request_id=_request_id(),
detail=decision.detail,
)
authorization = decision.to_dict()
if not allowed:
return JSONResponse(
{
"ok": False,
"error": "unauthorized",
"detail": (
"POST /api/v1/analytics/usage requires an authenticated "
"principal with record_analytics_usage execution enabled"
),
"authorization": authorization,
},
status_code=403,
)
usage_id = record_usage(
session_id=body.get("session_id"),
remote=body.get("remote", "dadeschools"),
org=body.get("org", ""),
repo=body.get("repo", ""),
project_id=body.get("project_id"),
role=body.get("role", "unknown"),
model=body.get("model", "unknown"),
issue_number=body.get("issue_number") or body.get("issue"),
pr_number=body.get("pr_number") or body.get("pr"),
stage=body.get("stage", "unknown"),
input_tokens=body.get("input_tokens"),
output_tokens=body.get("output_tokens"),
total_tokens=body.get("total_tokens"),
estimated_cost_usd=body.get("estimated_cost_usd"),
latency_ms=body.get("latency_ms"),
duration_ms=body.get("duration_ms"),
status=body.get("status", "success"),
metadata=body.get("metadata"),
)
return JSONResponse(
{"ok": True, "usage_id": usage_id, "authorization": authorization},
status_code=201,
)
async def method_not_allowed(request: Request, _exc: Exception) -> Response:
path = request.url.path
if path in _AUDIT_MUTATION_PATHS and request.method == "POST":
@@ -608,6 +722,10 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/runtime", runtime, methods=["GET"]),
Route("/api/runtime", api_runtime, methods=["GET"]),
Route("/api/v1/timeline", api_v1_timeline, methods=["GET"]),
Route("/analytics", analytics, methods=["GET"]),
Route("/api/analytics", api_v1_analytics, methods=["GET"]),
Route("/api/v1/analytics", api_v1_analytics, methods=["GET"]),
Route("/api/v1/analytics/usage", api_v1_analytics_ingest, methods=["POST"]),
Route("/audit", audit, methods=["GET", "POST"]),
Route("/api/audit", api_audit, methods=["GET", "POST"]),
Route("/worktrees", worktrees, methods=["GET"]),
+41
View File
@@ -236,6 +236,47 @@ _ACTION_SPECS: tuple[ConsoleAction, ...] = (
phase=3,
summary="Remove a remote feature branch.",
),
# #651 analytics ingest: local control-plane write, not a Gitea mutation.
# Phase 2 gated write so Phase 1 (ACTIVE_PHASE=1) fails closed on execution.
ConsoleAction(
action_id="record_analytics_usage",
task_key="record_analytics_usage",
action_class=CLASS_WRITE,
minimum_role=OPERATOR,
requires_confirmation=True,
dual_control=False,
break_glass=False,
phase=2,
summary="Ingest a model-usage / latency analytics event into the control-plane DB.",
),
# #642: sanctioned daemon lifecycle. These exist so operators have an
# audited path off `pkill -f mcp_server.py` (#630). Restart drops every
# in-flight request on a namespace, so it carries the same dual-control and
# break-glass weight as a merge; reload drains first and is privileged but
# not destructive. Neither ever exposes a raw kill: execution is handed to
# a host supervisor by ``webui.sanctioned_restart``.
ConsoleAction(
action_id="system.reload_namespace",
task_key="reload_namespace",
action_class=CLASS_PRIVILEGED,
minimum_role=CONTROLLER,
requires_confirmation=True,
dual_control=False,
break_glass=False,
phase=2,
summary="Gracefully reload one MCP namespace via the host supervisor.",
),
ConsoleAction(
action_id="system.restart_namespace",
task_key="restart_namespace",
action_class=CLASS_DESTRUCTIVE,
minimum_role=ADMIN,
requires_confirmation=True,
dual_control=True,
break_glass=True,
phase=2,
summary="Restart one MCP namespace via the host supervisor.",
),
)
ACTIONS: dict[str, ConsoleAction] = {a.action_id: a for a in _ACTION_SPECS}
+13
View File
@@ -110,6 +110,8 @@ def _format_target(action_id: str, params: dict[str, Any]) -> str:
)
if action_id == "create_issue":
return f"issue {params.get('title', '?')!r}"
if action_id in {"system.restart_namespace", "system.reload_namespace"}:
return f"MCP namespace {params.get('namespace', '?')!r}"
return "unspecified"
@@ -165,6 +167,17 @@ def build_action_registry() -> ActionRegistry:
"gitea_create_issue_comment", "Post a PR review thread comment."),
("close_pr", "Close PR", "close_pr", "gitea_edit_pr",
"Close a pull request without merge."),
# #642: the sanctioned replacement for the forbidden manual daemon-kill
# recovery path (#630). The "tool" is a host supervisor hook, not an MCP
# call — the console never signals a process. Preview and gating live in
# ``webui.sanctioned_restart``; these stay disabled like every other
# registry entry.
("system.reload_namespace", "Reload MCP namespace", "reload_namespace",
"host.supervisor_reload",
"Gracefully reload one MCP namespace via the host supervisor."),
("system.restart_namespace", "Restart MCP namespace",
"restart_namespace", "host.supervisor_restart",
"Restart one MCP namespace via the host supervisor."),
)
actions = tuple(
GatedAction(
+1
View File
@@ -65,6 +65,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
)),
NavGroup("Insights", (
NavItem("/insights", "Insights", "stub"),
NavItem("/analytics", "Analytics"),
NavItem("/audit", "Audit"),
)),
)
+613
View File
@@ -0,0 +1,613 @@
"""Sanctioned MCP restart and graceful reload controls (#642, Phase 2).
Sessions have historically recovered MCP connectivity by killing the host
daemon (``pkill -f mcp_server.py``, #630). That path stays forbidden: it kills
every namespace on the host, contaminates the surviving session, and leaves no
audit trail. This module is the sanctioned replacement.
A restart is modelled as a *gated action*, never as a command:
1. **Capability** — the console action resolves through ``console_authz``
against ``task_capability_map``, so the console cannot invent an authority
the MCP layer does not already define.
2. **Preview** — :func:`build_restart_preview` renders a mutation ledger and
the exact confirmation phrase. It never returns a shell command.
3. **Confirmation** — the operator echoes a phrase naming the exact namespace
and mode. A phrase for one namespace never authorizes another.
4. **Operator authorization** — host daemon maintenance is authorized out of
band through the environment (#630, and #710 finding F1: a worker session
cannot set an env var for an already-running daemon, so this cannot be
self-asserted the way a tool argument could).
5. **Execution** — :func:`execute_restart` never spawns a process. Once every
gate passes it hands the request to the configured host-managed restart
hook; with no hook configured it fails closed.
6. **Health recheck** — :func:`verify_post_restart_health` requires live
client-namespace probe evidence before any post-restart clean claim.
Manual ``pkill`` remains forbidden and is classified as contamination by
:func:`classify_restart_command`, which blocks clean claims (#630 AC3).
This module performs no I/O beyond reading its own environment configuration,
imports no MCP client, and holds no credential.
"""
from __future__ import annotations
import os
from dataclasses import asdict, dataclass
from typing import Any
import mcp_namespace_health
import runtime_recovery_guard
from webui import console_audit, console_authz
# --- Operations -------------------------------------------------------------
MODE_RESTART = "restart"
MODE_RELOAD = "reload"
MODES: tuple[str, ...] = (MODE_RESTART, MODE_RELOAD)
ACTION_RESTART_NAMESPACE = "system.restart_namespace"
ACTION_RELOAD_NAMESPACE = "system.reload_namespace"
ACTION_FOR_MODE: dict[str, str] = {
MODE_RESTART: ACTION_RESTART_NAMESPACE,
MODE_RELOAD: ACTION_RELOAD_NAMESPACE,
}
# Namespaces the console may target. An unlisted name fails closed rather than
# being passed through to a host hook.
KNOWN_NAMESPACES: tuple[str, ...] = tuple(
sorted(
set(mcp_namespace_health.DEFAULT_NAMESPACES)
| {"gitea-author", "gitea-reviewer", "gitea-merger",
"gitea-reconciler", "gitea-controller"}
)
)
# Scope tokens that would mean "everything at once". Explicit non-goal: the
# console never offers a fleet-wide restart, because that is the blast radius
# `pkill -f mcp_server.py` already had.
_FLEET_TOKENS = frozenset({"*", "all", "fleet", "any", ""})
# --- Environment configuration ----------------------------------------------
# Read server-side only; the value is an opaque host hook reference (e.g. a
# launchd label), never a command line, and is never rendered to a client.
RESTART_HOOK_ENV = "GITEA_SANCTIONED_RESTART_HOOK"
# --- Reason codes -----------------------------------------------------------
DENY_UNKNOWN_MODE = "unknown_mode"
DENY_UNKNOWN_NAMESPACE = "unknown_namespace"
DENY_FLEET_SCOPE = "fleet_scope_not_permitted"
DENY_UNAUTHORIZED = "unauthorized"
DENY_CONFIRMATION_MISSING = "confirmation_required"
DENY_CONFIRMATION_MISMATCH = "confirmation_mismatch"
DENY_OPERATOR_AUTHORIZATION = "operator_authorization_missing"
DENY_HOOK_NOT_CONFIGURED = "restart_hook_not_configured"
DENY_CONTAMINATED_RUNTIME = "contaminated_runtime"
ALLOW_HOST_ACTION_REQUIRED = "host_action_required"
# Post-restart verification outcomes.
HEALTH_CLEAN = "clean"
HEALTH_UNPROVEN = "unproven"
HEALTH_UNHEALTHY = "unhealthy"
def _clean(value: Any) -> str:
return str(value or "").strip()
# --- Mutation ledger --------------------------------------------------------
@dataclass(frozen=True)
class RestartLedgerEntry:
"""One planned step, shown before anything is asked of the host."""
sequence: int
step: str
summary: str
executes_process_kill: bool = False
def _mutation_ledger(namespace: str, mode: str) -> tuple[RestartLedgerEntry, ...]:
if mode == MODE_RELOAD:
middle = RestartLedgerEntry(
sequence=2,
step="host_graceful_reload",
summary=(
f"Ask the configured host supervisor to reload {namespace} "
"in place, draining in-flight requests. The console does not "
"signal the process itself."
),
)
else:
middle = RestartLedgerEntry(
sequence=2,
step="host_restart_hook",
summary=(
f"Ask the configured host supervisor to restart {namespace}. "
"The console never sends a signal and never runs a kill."
),
)
return (
RestartLedgerEntry(
sequence=1,
step="quiesce",
summary=(
f"Stop admitting new gated mutations for {namespace} and "
"record the intent before anything restarts."
),
),
middle,
RestartLedgerEntry(
sequence=3,
step="health_recheck",
summary=(
f"Re-probe {namespace} through the live client namespace and "
"prove the required tool is callable again."
),
),
RestartLedgerEntry(
sequence=4,
step="audit",
summary=(
"Append actor, target namespace, mode, and result to the "
"console audit log."
),
),
)
# --- Confirmation -----------------------------------------------------------
def confirmation_phrase(namespace: str, mode: str) -> str:
"""Exact phrase an operator must echo, naming the namespace and mode.
Binding the namespace into the phrase is the point: a confirmation typed
for ``gitea-author`` cannot be replayed against ``gitea-merger``.
"""
return f"{_clean(mode)} {_clean(namespace)}"
def confirmation_matches(
namespace: str, mode: str, confirmation: str | None
) -> bool:
"""Compare *confirmation* to the required phrase (exact, whitespace-trimmed)."""
return _clean(confirmation) == confirmation_phrase(namespace, mode)
# --- Scope validation -------------------------------------------------------
def _validate_scope(namespace: str, mode: str) -> tuple[str, str] | None:
"""Return ``(reason_code, detail)`` when the scope is refused."""
ns = _clean(namespace)
md = _clean(mode)
if md not in MODES:
return (
DENY_UNKNOWN_MODE,
f"Mode {md!r} is not one of {', '.join(MODES)}.",
)
if ns.lower() in _FLEET_TOKENS:
return (
DENY_FLEET_SCOPE,
(
"Fleet-wide restart is an explicit non-goal: it reproduces the "
"blast radius of `pkill -f mcp_server.py` (#630). Restart one "
"namespace at a time."
),
)
if ns not in KNOWN_NAMESPACES:
return (
DENY_UNKNOWN_NAMESPACE,
f"Namespace {ns!r} is not a known MCP namespace.",
)
return None
# --- Host hook --------------------------------------------------------------
def restart_hook(env: dict[str, str] | None = None) -> dict[str, Any]:
"""Report the configured host-managed restart hook.
The hook is a reference the *host* resolves (a supervisor label), not a
command this process runs. ``configured=False`` fails restart closed.
"""
source = env if env is not None else os.environ
reference = _clean(source.get(RESTART_HOOK_ENV))
return {
"configured": bool(reference),
"reference": reference or None,
"source": RESTART_HOOK_ENV if reference else None,
"self_assertable": False,
"console_executes_process": False,
}
# --- Preview ----------------------------------------------------------------
def build_restart_preview(
namespace: str,
mode: str = MODE_RESTART,
*,
principal: console_authz.Principal | None = None,
env: dict[str, str] | None = None,
) -> dict[str, Any]:
"""Render the dry-run preview for a restart/reload request.
Read-only: no authorization is granted, no host is contacted, and the
result never contains a shell command.
"""
ns = _clean(namespace)
md = _clean(mode)
action_id = ACTION_FOR_MODE.get(md, ACTION_RESTART_NAMESPACE)
action = console_authz.get_action(action_id)
decision = console_authz.authorize(action_id, principal)
scope_error = _validate_scope(ns, md)
hook = restart_hook(env)
operator = runtime_recovery_guard.operator_authorization(env)
return {
"action_id": action_id,
"namespace": ns,
"mode": md,
"scope_valid": scope_error is None,
"scope_reason_code": scope_error[0] if scope_error else None,
"scope_detail": scope_error[1] if scope_error else None,
"required_role": action.minimum_role if action else None,
"required_permission": action.mcp_permission if action else None,
"action_class": action.action_class if action else None,
"dual_control": action.dual_control if action else True,
"break_glass": action.break_glass if action else True,
"requires_confirmation": True,
"confirmation_phrase": confirmation_phrase(ns, md),
"mutation_ledger": [asdict(entry) for entry in _mutation_ledger(ns, md)],
"authorization": decision.to_dict(),
"operator_authorization": operator,
"restart_hook": hook,
"execution_enabled": False,
"raw_process_kill_exposed": False,
"known_namespaces": list(KNOWN_NAMESPACES),
"post_restart_verification_required": True,
}
# --- Gate -------------------------------------------------------------------
def assess_restart_request(
namespace: str,
mode: str = MODE_RESTART,
*,
principal: console_authz.Principal | None = None,
confirmation: str | None = None,
contamination_marker: dict[str, Any] | None = None,
env: dict[str, str] | None = None,
) -> dict[str, Any]:
"""Decide whether a restart request may proceed to the host hook.
Every gate must pass. The first failure wins and is reported with a stable
reason code; a pass never means "restarted", only "may be handed to the
configured host hook".
"""
ns = _clean(namespace)
md = _clean(mode)
action_id = ACTION_FOR_MODE.get(md, ACTION_RESTART_NAMESPACE)
preview = build_restart_preview(ns, md, principal=principal, env=env)
def refuse(reason_code: str, detail: str) -> dict[str, Any]:
return {
"allowed": False,
"gates_passed": False,
"reason_code": reason_code,
"detail": detail,
"action_id": action_id,
"namespace": ns,
"mode": md,
"preview": preview,
"execution_enabled": False,
}
scope_error = _validate_scope(ns, md)
if scope_error is not None:
return refuse(*scope_error)
# Authority is checked as an authorization decision, not an execution
# grant. ``for_execution=True`` asks "may the console perform this write?",
# and the answer here is permanently no: step 2 of the ledger is a request
# to the host supervisor, so the console's Phase 2 execution gate is not
# the relevant gate. Every branch below keeps ``execution_enabled`` False
# and :func:`execute_restart` never touches a process.
decision = console_authz.authorize(action_id, principal)
if not decision.allowed:
return refuse(DENY_UNAUTHORIZED, decision.detail)
if not _clean(confirmation):
return refuse(
DENY_CONFIRMATION_MISSING,
(
"Type the confirmation phrase "
f"{preview['confirmation_phrase']!r} to proceed."
),
)
if not confirmation_matches(ns, md, confirmation):
return refuse(
DENY_CONFIRMATION_MISMATCH,
(
"Confirmation does not name this namespace and mode; expected "
f"{preview['confirmation_phrase']!r}."
),
)
operator = preview["operator_authorization"]
if not operator["authorized"]:
return refuse(
DENY_OPERATOR_AUTHORIZATION,
(
"Host daemon maintenance requires out-of-band operator "
"authorization via "
f"{runtime_recovery_guard.OPERATOR_AUTHORIZATION_ENV}."
),
)
# #630's task-scoped gate deliberately lets a contaminated worker keep
# commenting and handing off. Restart is stricter and unconditional: a
# runtime already contaminated by a manual kill must be reconciled before
# it is restarted again, or the restart just launders the contamination.
if contamination_marker and not contamination_marker.get(
"cleared_by_reconciler"
):
return refuse(
DENY_CONTAMINATED_RUNTIME,
(
"A live contamination marker is present; clear it through the "
"reconciler path before restarting."
),
)
hook = preview["restart_hook"]
if not hook["configured"]:
return refuse(
DENY_HOOK_NOT_CONFIGURED,
(
"No host-managed restart hook is configured "
f"({RESTART_HOOK_ENV}). The console will not fall back to a "
"process kill."
),
)
return {
"allowed": True,
"gates_passed": True,
"reason_code": ALLOW_HOST_ACTION_REQUIRED,
"detail": (
"Every gate passed. The restart must be performed by the "
"configured host supervisor; the console does not signal the "
"process."
),
"action_id": action_id,
"namespace": ns,
"mode": md,
"preview": preview,
"execution_enabled": False,
"console_executes": False,
"console_active_phase": console_authz.ACTIVE_PHASE,
}
# --- Execution --------------------------------------------------------------
def execute_restart(
namespace: str,
mode: str = MODE_RESTART,
*,
principal: console_authz.Principal | None = None,
confirmation: str | None = None,
contamination_marker: dict[str, Any] | None = None,
env: dict[str, str] | None = None,
request_id: str | None = None,
session_id: str | None = None,
) -> dict[str, Any]:
"""Run every gate, audit the outcome, and hand off to the host.
This function never spawns a process, never sends a signal, and never
builds a command line. ``success`` is False in both directions: a refused
request is refused, and an authorized request still requires the host
supervisor to act.
"""
assessment = assess_restart_request(
namespace,
mode,
principal=principal,
confirmation=confirmation,
contamination_marker=contamination_marker,
env=env,
)
action_id = assessment["action_id"]
allowed = assessment["allowed"]
audit = console_audit.record_event(
action_id=action_id,
result=(
console_audit.RESULT_ALLOWED if allowed
else console_audit.RESULT_DENIED
),
principal=principal,
target={"namespace": assessment["namespace"], "mode": assessment["mode"]},
reason_code=assessment["reason_code"],
detail=assessment["detail"],
request_id=request_id,
session_id=session_id,
metadata={
"gates_passed": assessment["gates_passed"],
"process_kill_executed": False,
"post_restart_verification_required": True,
},
)
return {
"success": False,
"outcome": (
ALLOW_HOST_ACTION_REQUIRED if allowed else assessment["reason_code"]
),
"allowed": allowed,
"detail": assessment["detail"],
"namespace": assessment["namespace"],
"mode": assessment["mode"],
"action_id": action_id,
"process_kill_executed": False,
"host_hook": assessment["preview"]["restart_hook"],
"next_action": (
"Have the host supervisor perform the restart, then call "
"verify_post_restart_health with live client-namespace evidence "
"before claiming a clean session."
if allowed
else assessment["detail"]
),
"assessment": assessment,
"audit": audit,
}
# --- Contamination classification -------------------------------------------
def classify_restart_command(
command: str | None,
*,
mcp_pids: list[Any] | tuple[Any, ...] | None = None,
session_id: str | None = None,
remote: str | None = None,
role: str | None = None,
) -> dict[str, Any]:
"""Classify an operator-proposed recovery command (#630 AC2).
A manual ``pkill``/``kill`` of the MCP daemon is contamination, not a
restart. When contaminating, a durable marker is returned so downstream
gated mutations and clean claims fail closed.
"""
classification = runtime_recovery_guard.classify_recovery_command(
command, mcp_pids=mcp_pids
)
contaminating = bool(classification.get("contamination"))
marker = None
if contaminating:
marker = runtime_recovery_guard.build_contamination_record(
reason_class=(
classification.get("reason_class")
or runtime_recovery_guard.REASON_MANUAL_DAEMON_KILL
),
command_redacted=classification.get("redacted_command"),
session_id=session_id,
remote=remote,
role=role,
detail=(
"Manual daemon kill is forbidden; use the sanctioned "
f"{ACTION_RESTART_NAMESPACE} gated action instead."
),
)
return {
"contamination": contaminating,
"sanctioned": not contaminating and not classification.get("process_kill"),
"clean_claim_allowed": not contaminating,
"reason_class": classification.get("reason_class"),
"redacted_command": classification.get("redacted_command"),
"classification": classification,
"contamination_marker": marker,
"sanctioned_alternative": ACTION_RESTART_NAMESPACE,
}
# --- Post-restart health verification ---------------------------------------
def verify_post_restart_health(
namespace: str,
*,
probe_result: dict[str, Any] | None = None,
probe_source: str | None = None,
registered_tools: list[str] | tuple[str, ...] | None = None,
required_tool: str | None = None,
profile: str | None = None,
) -> dict[str, Any]:
"""Require live proof a namespace is callable before any clean claim (AC3).
Static registration is not proof and neither is an offline subprocess
probe: only ``probe_source=client_namespace`` evidence can clear a
post-restart session for mutations.
"""
ns = _clean(namespace)
health = mcp_namespace_health.classify_namespace_probe(
ns,
required_tool=required_tool,
registered_tools=registered_tools,
probe_result=probe_result,
profile=profile,
probe_source=probe_source,
)
healthy = bool(health.get("healthy"))
proven = bool(health.get("ide_namespace_proven"))
if healthy and proven:
status = HEALTH_CLEAN
elif healthy:
status = HEALTH_UNPROVEN
else:
status = HEALTH_UNHEALTHY
reasons = list(health.get("reasons") or [])
if status == HEALTH_UNPROVEN:
reasons.append(
"Namespace reported healthy without live client-namespace "
"evidence; a post-restart clean claim requires "
f"probe_source={mcp_namespace_health.PROBE_SOURCE_CLIENT!r}."
)
return {
"namespace": ns,
"status": status,
"healthy": healthy,
"ide_namespace_proven": proven,
"clean_claim_allowed": status == HEALTH_CLEAN,
"mutations_allowed": status == HEALTH_CLEAN,
"reasons": reasons,
"health": health,
}
# --- Policy surface ---------------------------------------------------------
def restart_policy() -> dict[str, Any]:
"""Machine-readable description of the sanctioned restart contract."""
return {
"policy_version": 1,
"modes": list(MODES),
"actions": [ACTION_RESTART_NAMESPACE, ACTION_RELOAD_NAMESPACE],
"known_namespaces": list(KNOWN_NAMESPACES),
"fleet_scope_permitted": False,
"console_executes_process_kill": False,
"raw_kill_exposed": False,
"requires_confirmation": True,
"confirmation_binds_namespace": True,
"operator_authorization_env": (
runtime_recovery_guard.OPERATOR_AUTHORIZATION_ENV
),
"restart_hook_env": RESTART_HOOK_ENV,
"manual_kill_classified_as": runtime_recovery_guard.CONTAMINATION_KIND,
"post_restart_clean_claim_requires": (
mcp_namespace_health.PROBE_SOURCE_CLIENT
),
"audit_required": True,
"silent_auto_restart_permitted": False,
}