Compare commits

..
Author SHA1 Message Date
sysadminandGrok 4.5 461e1dac78 feat(mcp): enforce scoped recovery playbook before full restart (Closes #669)
Add recovery_playbook.py with the narrow-to-broad recovery ladder, symptom
routing, attempt-log helpers, and escalation metrics. Wire the attempt-log
gate into restart_coordinator so rolling/full/host restarts require prior
insufficient narrower attempts (or break-glass). Document the ladder and
update gitea_request_mcp_restart for prior_recovery_attempts_json.

Co-Authored-By: Grok 4.5 <[email protected]>
2026-07-25 17:35:11 -04:00
13 changed files with 1018 additions and 1428 deletions
+94
View File
@@ -0,0 +1,94 @@
# MCP scoped recovery playbook (#669)
**Parent:** [#655](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/655)
**Vision / roadmap:** [#652](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/652) · [#653](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/653)
**Class matrix:** [#663](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/663) · `docs/mcp-restart-classes.md`
**Coordinator:** [#658](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/658) · `restart_coordinator.py`
**Audit lineage:** [#665](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/665)
## Decision
Full-server MCP reset is a **last resort**. Prefer the narrowest recovery that
can clear the symptom. The coordinator **refuses** `rolling_mcp_restart`,
`full_mcp_restart`, and `host_restart` unless:
1. The inventory carries a prior **attempt log** of at least one *insufficient*
narrower recovery, **or**
2. **Break-glass** is authorized
(`request_break_glass` + `GITEA_BREAKGLASS_RESTART_AUTHORIZATION`).
Break-glass still never bypasses the #663 class matrix (role/permission).
## Ladder (narrow → broad)
| Rank | Action | Self-service | Implementation / delegation |
|---:|---|---|---|
| 0 | `client_reconnect` | yes | Host auto-reconnect / client reconnect · [#584](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/584) · `docs/mcp-namespace-eof-recovery.md` |
| 1 | `capability_refresh` | yes | `gitea_resolve_task_capability` + `gitea_whoami` · [#610](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/610) · [#685](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/685) |
| 2 | `session_reconnect` | yes | Runtime rebind + explicit `worktree_path` · [#543](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/543) · [#618](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/618) |
| 3 | `configuration_reload` | no | Class `configuration_reload` · console reload · [#642](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/642) |
| 4 | `lease_recovery` | no | Lock/lease recovery paths · [#702](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/702) · [#753](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/753) · [#790](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/790) |
| 5 | `worker_restart` | no | Class `worker_restart` · #663 |
| 6 | `role_runtime_restart` | no | Class `role_runtime_restart` · console restart · #642/#663 |
| 7 | `connector_restart` | no | Class `connector_restart` · #663 |
| 8 | `rolling_mcp_restart` | no | Class `rolling_mcp_restart` · design [#668](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/668) · **attempt log required** |
| 9 | `full_mcp_restart` | no | Class `full_mcp_restart` · **attempt log required** |
| 10 | `host_restart` | no | Class `host_restart` · **attempt log required** |
Machine-readable source of truth: `recovery_playbook.RECOVERY_LADDER` and
`recovery_playbook.ladder_document()`.
## Attempt log shape
Each prior attempt is a mapping:
```json
{
"action": "client_reconnect",
"outcome": "insufficient",
"reason": "transport still closed after IDE reconnect",
"actor": "prgs-controller-12345",
"recorded_at": "2026-07-25T21:00:00+00:00"
}
```
Outcomes that count toward escalation: `failed`, `insufficient`, `denied`,
`unresolved`, `timeout`, `error`.
Pass attempts into the coordinator via inventory
`prior_recovery_attempts` or the MCP tool argument
`prior_recovery_attempts_json` on `gitea_request_mcp_restart`.
Helper: `recovery_playbook.build_attempt_record(...)`.
## Symptom → first rung
`recovery_playbook.recommend_actions(symptoms=[...])` maps symptoms such as
`transport_eof`, `stale_capability`, `stale_lease`, `daemon_corrupt` to the
narrowest recommended action, then walks the ladder. Soft recommendations
never replace the hard gate on broad restarts.
## Enforcement points
1. **`recovery_playbook.assess_escalation`** — pure gate.
2. **`restart_coordinator.evaluate_restart_impact`** — when `restart_class` is
set (policy-enforced path), broad classes require the gate; report fields
`attempt_log_satisfied`, `playbook_escalation`, `break_glass`.
3. **`gitea_request_mcp_restart`** — accepts attempt JSON and env-authorized
break-glass; never restarts a process.
## Metrics
`recovery_playbook.recovery_metrics(attempts)` reports the fraction of
successful recoveries that avoided full/host restart
(`fraction_avoided_full_restart`).
## Non-goals
* HA multi-instance execution (#668 design only here).
* Normalizing `pkill` (#630 contamination stays forbidden).
* Silent mutation of leases or processes from the playbook itself.
## Manual process kills
Remain forbidden and contaminating (#630). The playbook never recommends them.
+15 -5
View File
@@ -95,12 +95,18 @@ gitea_request_mcp_restart(remote, host, org, repo,
target_session_id=None, target_role=None, target_session_id=None, target_role=None,
target_connector=None, target_connector=None,
drain_proof_json=None, drain_proof_json=None,
request_break_glass=False) request_break_glass=False,
prior_recovery_attempts_json=None)
``` ```
It **never restarts anything**: `apply_supported` is always `false` and It **never restarts anything**: `apply_supported` is always `false` and
`restart_performed` is always `false`. `restart_performed` is always `false`.
`prior_recovery_attempts_json` (#669) is an optional JSON array of prior
narrow recovery attempts. Rolling / full / host classes require at least one
*insufficient* narrower attempt (or authorized break-glass). See
`docs/mcp-recovery-playbook.md`.
### Dry-run versus apply ### Dry-run versus apply
| Call | Behavior | | Call | Behavior |
@@ -112,9 +118,10 @@ It **never restarts anything**: `apply_supported` is always `false` and
An apply requires **both** authorizations, and they are independent: An apply requires **both** authorizations, and they are independent:
1. **Restart-class authorization** (#663) — the requester's role and permissions 1. **Restart-class authorization** (#663 / #669) — the requester's role and
must allow the requested class, the class's approval requirement must be permissions must allow the requested class, the class's approval requirement
satisfied, and any target-scoped class must name its target. Failing any of must be satisfied, any target-scoped class must name its target, and broad
classes must satisfy the recovery-playbook attempt-log gate. Failing any of
these makes `allow_restart` `false`. these makes `allow_restart` `false`.
2. **Drain-proof gate** (#661) — a valid, unexpired, clean proof bound to the 2. **Drain-proof gate** (#661) — a valid, unexpired, clean proof bound to the
current impact fingerprint, or an authorized break-glass. current impact fingerprint, or an authorized break-glass.
@@ -127,7 +134,10 @@ the authorization that produced it.
### Break-glass ### Break-glass
Break-glass bypasses the **drain proof only** — never the restart-class matrix. Break-glass bypasses the **drain proof only** — never the restart-class matrix
(role/permission). Separately, authorized break-glass also satisfies the #669
attempt-log requirement for broad restarts (rolling/full/host), because that
gate is not a class-matrix permission check.
It is honoured solely when `request_break_glass` is set *and* the environment It is honoured solely when `request_break_glass` is set *and* the environment
carries `GITEA_BREAKGLASS_RESTART_AUTHORIZATION`; like operator override, the carries `GITEA_BREAKGLASS_RESTART_AUTHORIZATION`; like operator override, the
tool argument expresses caller intent and cannot be self-asserted by a worker tool argument expresses caller intent and cannot be self-asserted by a worker
+1 -54
View File
@@ -85,10 +85,7 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| `/inventory` | Phase 1 shell stub — unified inventory (backed by #636) | | `/inventory` | Phase 1 shell stub — unified inventory (backed by #636) |
| `/timeline` | Phase 1 shell stub — workflow event timeline | | `/timeline` | Phase 1 shell stub — workflow event timeline |
| `/policy` | Phase 1 shell stub — capability/role policy placeholder | | `/policy` | Phase 1 shell stub — capability/role policy placeholder |
| `/providers` | AI-provider connection status (#650) — declared registry only, no secrets | | `/insights` | Phase 1 shell stub — operational insights placeholder |
| `/api/v1/providers` | JSON provider connections; `502` when the registry cannot be loaded |
| `/insights` | Evidence-backed operational insights (#650) — advisory only |
| `/api/v1/insights` | JSON insights export with evidence refs and source availability |
Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for `read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
@@ -332,56 +329,6 @@ Honesty rules specific to this view:
The write-time redactor is a narrow denylist and is not relied on. The field The write-time redactor is a narrow denylist and is not relied on. The field
itself is kept — it is the `#630` evidence naming which daemon was killed. itself is kept — it is the `#630` evidence naming which daemon was killed.
## AI providers and operational insights (#650)
`/providers` and `/insights` are the Phase 4 **advisory** surfaces for AI-provider
connections and evidence-backed operational findings. They never expose API keys,
never mutate Gitea, and never authorize review, merge, or close.
### Provider connections (`/providers`)
Status is taken from the **worker registry** declaration (`webui/data/workers.registry.json`
or `WEBUI_WORKER_REGISTRY`):
| Field | Meaning |
|-------|---------|
| `connection_status` | `declared_available` or `declared_unavailable` from the registry `available` flag |
| `models` | Declared model list only (not a live vendor enumeration) |
| `worker_count` / `enabled_worker_count` | How many worker instances name this provider |
| `secrets_exposed` | Always `false` — credentials are never loaded |
Live executable health is **not** probed here (that belongs to the provider adapter
framework). The page states this probe limit explicitly so a green badge is not
misread as a process heartbeat.
`GET /api/v1/providers` returns the same model (`schema_version: 1`). It answers
`502` when the registry cannot be loaded so consumers cannot treat a fail-closed
payload as “no providers configured”.
### Operational insights (`/insights`)
Insights are pure functions over durable console evidence:
| Kind | Evidence source |
|------|-----------------|
| `blocked_queue_pressure` | Traffic control blocked bucket (issue/PR numbers + reasons) |
| `controller_attention` | Traffic control `needs_controller` items |
| `stale_runtime_risk` | System-health stale_runtime / mutation_safe |
| `provider_without_workers` | Declared-available providers with zero workers |
| `analytics_failure_pressure` | Analytics events with failure status (when loaded) |
Rules:
* Every insight carries at least one evidence ref (`kind` + `ref` + `detail`).
Evidence-less insights are refused, not emitted.
* `advisory_only` is always true; `claims_action_completed` is always false.
* Missing sources appear under `sources_unavailable` — never as a silent empty
“all clear”.
* Titles and details pass through console redaction before display.
`GET /api/v1/insights` exports the same model. The HTML page always renders
interpretation limits so operators know these cards do not override workflow gates.
## Gitea issue/PR linkage (#645) ## Gitea issue/PR linkage (#645)
`/gitea` is the Phase 3 read-only linkage console: which PR carries which issue, `/gitea` is the Phase 3 read-only linkage console: which PR carries which issue,
+35 -9
View File
@@ -22575,8 +22575,9 @@ def gitea_request_mcp_restart(
target_connector: str | None = None, target_connector: str | None = None,
drain_proof_json: str | None = None, drain_proof_json: str | None = None,
request_break_glass: bool = False, request_break_glass: bool = False,
prior_recovery_attempts_json: str | None = None,
) -> dict: ) -> dict:
"""Evaluate a proposed MCP restart and return an impact preview (#658). """Evaluate a proposed MCP restart and return an impact preview (#658/#669).
Central restart coordinator: resolves the requested restart class, gathers Central restart coordinator: resolves the requested restart class, gathers
live control-plane state (sessions, live control-plane state (sessions,
@@ -22600,10 +22601,16 @@ def gitea_request_mcp_restart(
independent the drain gate proves the blast radius was drained and knows independent the drain gate proves the blast radius was drained and knows
nothing about whether this requester may request this class so a class the nothing about whether this requester may request this class so a class the
matrix denied never reports an authorized apply. Break-glass bypasses the matrix denied never reports an authorized apply. Break-glass bypasses the
drain proof only; it never bypasses the class matrix. ``apply_gate`` carries drain proof and, when env-authorized, the #669 attempt-log requirement for
broad restarts; it never bypasses the class matrix. ``apply_gate`` carries
``drain_gate_allow`` and ``restart_class_authorized`` so a denial is ``drain_gate_allow`` and ``restart_class_authorized`` so a denial is
attributable to the authorization that produced it. attributable to the authorization that produced it.
``prior_recovery_attempts_json`` (#669) is an optional JSON array of prior
narrow recovery attempts ``{action, outcome, reason, ...}``. Rolling / full
/ host restart classes require at least one *insufficient* narrower attempt
unless break-glass is authorized.
Operator override authority is read from the process environment Operator override authority is read from the process environment
(``GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION``), never self-asserted by (``GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION``), never self-asserted by
the requesting session: ``request_override`` only expresses caller intent the requesting session: ``request_override`` only expresses caller intent
@@ -22709,12 +22716,37 @@ def gitea_request_mcp_restart(
requester_role requester_role
) )
prior_recovery_attempts: list[dict] = []
if prior_recovery_attempts_json:
try:
parsed_attempts = json.loads(prior_recovery_attempts_json)
if isinstance(parsed_attempts, list):
prior_recovery_attempts = [
dict(a) for a in parsed_attempts if isinstance(a, dict)
]
else:
incomplete_reasons.append(
"prior_recovery_attempts_json must be a JSON array (#669)"
)
inventory_complete = False
except (ValueError, TypeError) as exc:
incomplete_reasons.append(
f"invalid prior_recovery_attempts_json: {_redact(str(exc))}"
)
inventory_complete = False
break_glass_authorized = bool(
(os.environ.get("GITEA_BREAKGLASS_RESTART_AUTHORIZATION") or "").strip()
)
break_glass = bool(request_break_glass and break_glass_authorized)
inventory = { inventory = {
"sessions": sessions, "sessions": sessions,
"leases": leases, "leases": leases,
"terminal_lock": terminal_lock, "terminal_lock": terminal_lock,
"inventory_complete": inventory_complete, "inventory_complete": inventory_complete,
"incomplete_reasons": incomplete_reasons, "incomplete_reasons": incomplete_reasons,
"prior_recovery_attempts": prior_recovery_attempts,
} }
report = restart_coordinator.evaluate_restart_impact( report = restart_coordinator.evaluate_restart_impact(
@@ -22730,6 +22762,7 @@ def gitea_request_mcp_restart(
target_session_id=target_session_id, target_session_id=target_session_id,
target_role=target_role, target_role=target_role,
target_connector=target_connector, target_connector=target_connector,
break_glass=break_glass,
) )
payload = report.as_dict() payload = report.as_dict()
@@ -22762,13 +22795,6 @@ def gitea_request_mcp_restart(
except (ValueError, TypeError) as exc: except (ValueError, TypeError) as exc:
proof_parse_error = f"invalid drain_proof_json: {_redact(str(exc))}" proof_parse_error = f"invalid drain_proof_json: {_redact(str(exc))}"
break_glass_authorized = bool(
(
os.environ.get("GITEA_BREAKGLASS_RESTART_AUTHORIZATION") or ""
).strip()
)
break_glass = bool(request_break_glass and break_glass_authorized)
expected_fp = drain_proof.impact_fingerprint(report.as_dict()) expected_fp = drain_proof.impact_fingerprint(report.as_dict())
gate = drain_proof.gate_apply_restart( gate = drain_proof.gate_apply_restart(
proof=proof_obj, proof=proof_obj,
+583
View File
@@ -0,0 +1,583 @@
"""Scoped MCP recovery playbook (#669).
Operational recovery must prefer the *narrowest* action that can fix the
symptom. Full MCP / host restarts are last-resort rungs on a documented
ladder; the coordinator refuses those rungs unless a prior attempt log
shows narrower recoveries already failed (or break-glass is authorized).
This module is pure classification and recommendation:
* No network, filesystem, or process I/O.
* Never restarts anything.
* Narrow recovery *execution* is delegated to existing tools/docs (linked
per rung) — the playbook records which rung to try next and whether
escalation to a broad restart is allowed.
Design lineage: umbrella #655, class matrix #663, coordinator #658,
auto-reconnect #584, stale-runtime #610, contamination #630, audit #665.
Vision #652 / roadmap #653.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from datetime import datetime, timezone
from enum import Enum
from typing import Any, Mapping, Sequence
PLAYBOOK_VERSION = "1.0.0-issue-669"
# Attempt outcomes that count as "tried and insufficient" for escalation.
INSUFFICIENT_OUTCOMES = frozenset(
{
"failed",
"insufficient",
"denied",
"unresolved",
"timeout",
"error",
}
)
# Break-glass / operator override still records that the ladder was skipped.
OUTCOME_BREAK_GLASS = "break_glass"
OUTCOME_SUCCESS = "success"
OUTCOME_SKIPPED = "skipped"
class RecoveryAction(str, Enum):
"""Ordered recovery ladder (narrow → broad)."""
CLIENT_RECONNECT = "client_reconnect"
CAPABILITY_REFRESH = "capability_refresh"
SESSION_RECONNECT = "session_reconnect"
CONFIGURATION_RELOAD = "configuration_reload"
LEASE_RECOVERY = "lease_recovery"
WORKER_RESTART = "worker_restart"
ROLE_RUNTIME_RESTART = "role_runtime_restart"
CONNECTOR_RESTART = "connector_restart"
ROLLING_MCP_RESTART = "rolling_mcp_restart"
FULL_MCP_RESTART = "full_mcp_restart"
HOST_RESTART = "host_restart"
# Classes that require a prior narrow-attempt log (unless break-glass).
BROAD_RESTART_ACTIONS: frozenset[RecoveryAction] = frozenset(
{
RecoveryAction.ROLLING_MCP_RESTART,
RecoveryAction.FULL_MCP_RESTART,
RecoveryAction.HOST_RESTART,
}
)
# Map #663 restart_class strings onto playbook actions.
RESTART_CLASS_TO_ACTION: dict[str, RecoveryAction] = {
"client_reconnect": RecoveryAction.CLIENT_RECONNECT,
"session_reconnect": RecoveryAction.SESSION_RECONNECT,
"configuration_reload": RecoveryAction.CONFIGURATION_RELOAD,
"worker_restart": RecoveryAction.WORKER_RESTART,
"role_runtime_restart": RecoveryAction.ROLE_RUNTIME_RESTART,
"connector_restart": RecoveryAction.CONNECTOR_RESTART,
"rolling_mcp_restart": RecoveryAction.ROLLING_MCP_RESTART,
"full_mcp_restart": RecoveryAction.FULL_MCP_RESTART,
"host_restart": RecoveryAction.HOST_RESTART,
}
@dataclass(frozen=True)
class RecoveryRung:
"""One rung on the recovery ladder."""
action: RecoveryAction
rank: int
summary: str
# Existing implementation or explicit delegation target.
implementation: str
issue_links: tuple[str, ...]
self_service: bool
# Restart-class permission when this rung is requested via coordinator.
restart_class: str | None = None
def as_dict(self) -> dict[str, Any]:
return {
"action": self.action.value,
"rank": self.rank,
"summary": self.summary,
"implementation": self.implementation,
"issue_links": list(self.issue_links),
"self_service": self.self_service,
"restart_class": self.restart_class,
}
# Canonical ladder. Rank 0 is narrowest.
RECOVERY_LADDER: tuple[RecoveryRung, ...] = (
RecoveryRung(
RecoveryAction.CLIENT_RECONNECT,
0,
"Reconnect the IDE/client MCP transport (EOF / transport flap).",
"Host auto-reconnect or explicit client reconnect; "
"docs/mcp-namespace-eof-recovery.md",
("#584", "#655"),
True,
"client_reconnect",
),
RecoveryRung(
RecoveryAction.CAPABILITY_REFRESH,
1,
"Re-resolve task capability and clear stale permission context.",
"Delegated: gitea_resolve_task_capability + gitea_whoami "
"(no process change).",
("#610", "#685", "#655"),
True,
None,
),
RecoveryRung(
RecoveryAction.SESSION_RECONNECT,
2,
"Rebind identity, workspace, and namespace for one session.",
"Delegated: gitea_get_runtime_context + explicit worktree_path "
"rebind (#618); docs/mcp-namespace-health.md",
("#543", "#618", "#655"),
True,
"session_reconnect",
),
RecoveryRung(
RecoveryAction.CONFIGURATION_RELOAD,
3,
"Gracefully reload configuration without replacing the daemon.",
"restart_coordinator class configuration_reload; console "
"system.reload_namespace (#642).",
("#642", "#663", "#655"),
False,
"configuration_reload",
),
RecoveryRung(
RecoveryAction.LEASE_RECOVERY,
4,
"Recover or rebind stale leases/locks without a process restart.",
"Delegated: issue lock recovery / lease lifecycle paths "
"(#702, #753, #790).",
("#702", "#753", "#790", "#655"),
False,
None,
),
RecoveryRung(
RecoveryAction.WORKER_RESTART,
5,
"Restart one worker after its own lease and mutation scope drains.",
"restart_coordinator class worker_restart (#663).",
("#663", "#655"),
False,
"worker_restart",
),
RecoveryRung(
RecoveryAction.ROLE_RUNTIME_RESTART,
6,
"Restart one role runtime and re-probe that namespace only.",
"restart_coordinator class role_runtime_restart; console "
"system.restart_namespace (#642).",
("#642", "#663", "#655"),
False,
"role_runtime_restart",
),
RecoveryRung(
RecoveryAction.CONNECTOR_RESTART,
7,
"Restart one connector while unrelated runtimes stay available.",
"restart_coordinator class connector_restart (#663).",
("#663", "#655"),
False,
"connector_restart",
),
RecoveryRung(
RecoveryAction.ROLLING_MCP_RESTART,
8,
"Drain/restart/verify one instance at a time (HA path).",
"restart_coordinator class rolling_mcp_restart; design #668.",
("#668", "#663", "#655"),
False,
"rolling_mcp_restart",
),
RecoveryRung(
RecoveryAction.FULL_MCP_RESTART,
9,
"Full stable-control MCP process restart after verified full drain.",
"restart_coordinator class full_mcp_restart; requires attempt log "
"unless break-glass (#669).",
("#658", "#661", "#663", "#669", "#655"),
False,
"full_mcp_restart",
),
RecoveryRung(
RecoveryAction.HOST_RESTART,
10,
"Host/infrastructure restart — broadest last-resort action.",
"restart_coordinator class host_restart; operator-owned.",
("#663", "#669", "#655"),
False,
"host_restart",
),
)
_LADDER_BY_ACTION: dict[RecoveryAction, RecoveryRung] = {
rung.action: rung for rung in RECOVERY_LADDER
}
# Symptom tokens → preferred first rung (decision tree, #663 lineage).
SYMPTOM_TO_FIRST_ACTION: dict[str, RecoveryAction] = {
"transport_eof": RecoveryAction.CLIENT_RECONNECT,
"client_closing_eof": RecoveryAction.CLIENT_RECONNECT,
"transport_flap": RecoveryAction.CLIENT_RECONNECT,
"namespace_disconnected": RecoveryAction.CLIENT_RECONNECT,
"stale_capability": RecoveryAction.CAPABILITY_REFRESH,
"permission_stale": RecoveryAction.CAPABILITY_REFRESH,
"runtime_reconnect_required": RecoveryAction.CAPABILITY_REFRESH,
"stale_runtime": RecoveryAction.SESSION_RECONNECT,
"worktree_unbound": RecoveryAction.SESSION_RECONNECT,
"namespace_unhealthy": RecoveryAction.SESSION_RECONNECT,
"config_drift": RecoveryAction.CONFIGURATION_RELOAD,
"profile_misbound": RecoveryAction.CONFIGURATION_RELOAD,
"stale_lease": RecoveryAction.LEASE_RECOVERY,
"dead_pid_lock": RecoveryAction.LEASE_RECOVERY,
"orphan_worktree": RecoveryAction.LEASE_RECOVERY,
"single_worker_stuck": RecoveryAction.WORKER_RESTART,
"role_runtime_dead": RecoveryAction.ROLE_RUNTIME_RESTART,
"connector_dead": RecoveryAction.CONNECTOR_RESTART,
"ha_instance_unhealthy": RecoveryAction.ROLLING_MCP_RESTART,
"daemon_corrupt": RecoveryAction.FULL_MCP_RESTART,
"full_process_deadlock": RecoveryAction.FULL_MCP_RESTART,
"host_unresponsive": RecoveryAction.HOST_RESTART,
}
def _utc_now() -> datetime:
return datetime.now(timezone.utc)
def resolve_action(value: RecoveryAction | str) -> RecoveryAction:
"""Resolve a recovery action or fail closed for unknown values."""
if isinstance(value, RecoveryAction):
return value
text = str(value or "").strip()
# Accept #663 restart_class aliases.
if text in RESTART_CLASS_TO_ACTION:
return RESTART_CLASS_TO_ACTION[text]
try:
return RecoveryAction(text)
except ValueError as exc:
raise ValueError(
f"unknown recovery action {value!r}; deny (fail closed, #669)"
) from exc
def ladder_rank(action: RecoveryAction | str) -> int:
resolved = resolve_action(action)
return _LADDER_BY_ACTION[resolved].rank
def rung_for(action: RecoveryAction | str) -> RecoveryRung:
return _LADDER_BY_ACTION[resolve_action(action)]
def normalize_attempt(raw: Mapping[str, Any]) -> dict[str, Any] | None:
"""Normalize one prior-recovery attempt record; return None if unusable."""
if not isinstance(raw, Mapping):
return None
action_raw = raw.get("action") or raw.get("recovery_action") or raw.get(
"restart_class"
)
if not action_raw:
return None
try:
action = resolve_action(str(action_raw))
except ValueError:
return None
outcome = str(
raw.get("outcome") or raw.get("status") or raw.get("result") or ""
).strip().lower()
if not outcome:
return None
recorded_at = raw.get("recorded_at") or raw.get("at") or raw.get("timestamp")
reason = str(raw.get("reason") or raw.get("detail") or "").strip()
actor = str(raw.get("actor") or raw.get("session_id") or "").strip()
return {
"action": action.value,
"outcome": outcome,
"reason": reason,
"actor": actor,
"recorded_at": recorded_at,
"rank": ladder_rank(action),
"raw": dict(raw),
}
def normalize_attempt_log(
attempts: Sequence[Mapping[str, Any]] | None,
) -> list[dict[str, Any]]:
"""Return usable attempt records in ladder order."""
out: list[dict[str, Any]] = []
for raw in attempts or ():
norm = normalize_attempt(raw)
if norm is not None:
out.append(norm)
out.sort(key=lambda a: (a["rank"], str(a.get("recorded_at") or "")))
return out
def narrower_insufficient_attempts(
attempts: Sequence[Mapping[str, Any]] | None,
*,
requested: RecoveryAction | str,
) -> list[dict[str, Any]]:
"""Return prior attempts narrower than *requested* that were insufficient."""
target_rank = ladder_rank(requested)
usable = []
for attempt in normalize_attempt_log(attempts):
if attempt["rank"] >= target_rank:
continue
if attempt["outcome"] in INSUFFICIENT_OUTCOMES:
usable.append(attempt)
return usable
@dataclass(frozen=True)
class EscalationAssessment:
"""Whether a requested broad recovery may proceed given the attempt log."""
requested_action: str
allowed: bool
require_attempt_log: bool
break_glass: bool
reasons: list[str] = field(default_factory=list)
qualifying_attempts: list[dict[str, Any]] = field(default_factory=list)
recommended_next: list[dict[str, Any]] = field(default_factory=list)
playbook_version: str = PLAYBOOK_VERSION
def as_dict(self) -> dict[str, Any]:
return {
"playbook_version": self.playbook_version,
"requested_action": self.requested_action,
"allowed": self.allowed,
"require_attempt_log": self.require_attempt_log,
"break_glass": self.break_glass,
"reasons": list(self.reasons),
"qualifying_attempts": list(self.qualifying_attempts),
"recommended_next": list(self.recommended_next),
}
def assess_escalation(
requested: RecoveryAction | str,
*,
prior_recovery_attempts: Sequence[Mapping[str, Any]] | None = None,
break_glass: bool = False,
) -> EscalationAssessment:
"""Gate broad restarts on a prior narrow-attempt log (#669 AC3).
Narrow / mid-ladder actions do not require a prior attempt log.
``full_mcp_restart``, ``host_restart``, and ``rolling_mcp_restart``
require at least one *insufficient* narrower attempt unless
``break_glass`` is true.
"""
action = resolve_action(requested)
require_log = action in BROAD_RESTART_ACTIONS
reasons: list[str] = []
qualifying = narrower_insufficient_attempts(
prior_recovery_attempts, requested=action
)
if not require_log:
return EscalationAssessment(
requested_action=action.value,
allowed=True,
require_attempt_log=False,
break_glass=bool(break_glass),
reasons=["narrow recovery; attempt log not required"],
qualifying_attempts=qualifying,
recommended_next=[],
)
if break_glass:
return EscalationAssessment(
requested_action=action.value,
allowed=True,
require_attempt_log=True,
break_glass=True,
reasons=[
"break-glass authorized; broad restart permitted without "
"narrow-attempt log (#669)"
],
qualifying_attempts=qualifying,
recommended_next=[],
)
if qualifying:
return EscalationAssessment(
requested_action=action.value,
allowed=True,
require_attempt_log=True,
break_glass=False,
reasons=[
f"{len(qualifying)} narrower recovery attempt(s) recorded as "
"insufficient; escalation permitted"
],
qualifying_attempts=qualifying,
recommended_next=[],
)
# Deny: recommend the next untried narrow rung(s).
recommended = recommend_actions(
symptoms=(),
prior_recovery_attempts=prior_recovery_attempts,
max_actions=3,
)
reasons.append(
f"{action.value} requires a prior attempt log of insufficient "
"narrower recoveries (or break-glass); none found — deny (fail "
"closed, #669)"
)
return EscalationAssessment(
requested_action=action.value,
allowed=False,
require_attempt_log=True,
break_glass=False,
reasons=reasons,
qualifying_attempts=[],
recommended_next=recommended.get("recommended_actions") or [],
)
def recommend_actions(
*,
symptoms: Sequence[str] = (),
prior_recovery_attempts: Sequence[Mapping[str, Any]] | None = None,
max_actions: int = 5,
) -> dict[str, Any]:
"""Return ordered recommended recovery actions for the given symptoms.
Soft mode (rollout): recommendations only — callers decide whether to
hard-gate. Hard mode for broad restarts is :func:`assess_escalation`.
"""
attempted_success = {
a["action"]
for a in normalize_attempt_log(prior_recovery_attempts)
if a["outcome"] == OUTCOME_SUCCESS
}
attempted_any = {
a["action"] for a in normalize_attempt_log(prior_recovery_attempts)
}
first_actions: list[RecoveryAction] = []
for symptom in symptoms:
key = str(symptom or "").strip().lower().replace(" ", "_").replace("-", "_")
mapped = SYMPTOM_TO_FIRST_ACTION.get(key)
if mapped is not None and mapped not in first_actions:
first_actions.append(mapped)
# Default entry: client reconnect then walk the ladder.
if not first_actions:
first_actions = [RecoveryAction.CLIENT_RECONNECT]
recommended: list[dict[str, Any]] = []
seen: set[str] = set()
min_rank = min(ladder_rank(a) for a in first_actions)
for rung in RECOVERY_LADDER:
if rung.rank < min_rank:
continue
if rung.action.value in attempted_success:
continue
if rung.action.value in seen:
continue
# Prefer rungs not yet attempted; still list previously-failed ones
# only if nothing else remains.
entry = rung.as_dict()
entry["already_attempted"] = rung.action.value in attempted_any
recommended.append(entry)
seen.add(rung.action.value)
if len(recommended) >= max(1, int(max_actions)):
break
return {
"playbook_version": PLAYBOOK_VERSION,
"symptoms": [str(s) for s in symptoms],
"recommended_actions": recommended,
"ladder": [r.as_dict() for r in RECOVERY_LADDER],
"read_only": True,
"hard_gate_note": (
"Broad restarts (rolling/full/host) still require "
"assess_escalation / coordinator attempt-log enforcement."
),
}
def build_attempt_record(
action: RecoveryAction | str,
*,
outcome: str,
reason: str = "",
actor: str = "",
recorded_at: str | None = None,
extra: Mapping[str, Any] | None = None,
) -> dict[str, Any]:
"""Build a durable-shaped attempt log entry for inventory/audit (#665)."""
resolved = resolve_action(action)
record = {
"action": resolved.value,
"outcome": str(outcome or "").strip().lower(),
"reason": str(reason or "").strip(),
"actor": str(actor or "").strip(),
"recorded_at": recorded_at or _utc_now().isoformat(),
"rank": ladder_rank(resolved),
"playbook_version": PLAYBOOK_VERSION,
}
if extra:
record["extra"] = dict(extra)
return record
def recovery_metrics(
attempts: Sequence[Mapping[str, Any]] | None,
) -> dict[str, Any]:
"""Compute the fraction of recoveries that avoided full/host restart.
A recovery *episode* is approximated as one attempt with
``outcome=success``. Successes on non-broad rungs count as avoided full
restart; successes on full/host count as full-restart recoveries.
"""
norms = normalize_attempt_log(attempts)
successes = [a for a in norms if a["outcome"] == OUTCOME_SUCCESS]
broad_success = [
a
for a in successes
if resolve_action(a["action"])
in {RecoveryAction.FULL_MCP_RESTART, RecoveryAction.HOST_RESTART}
]
avoided = [a for a in successes if a not in broad_success]
total = len(successes)
fraction_avoided = (len(avoided) / total) if total else None
return {
"playbook_version": PLAYBOOK_VERSION,
"attempts_total": len(norms),
"successes_total": total,
"successes_avoided_full_restart": len(avoided),
"successes_full_or_host_restart": len(broad_success),
"fraction_avoided_full_restart": fraction_avoided,
"insufficient_attempts": sum(
1 for a in norms if a["outcome"] in INSUFFICIENT_OUTCOMES
),
}
def ladder_document() -> dict[str, Any]:
"""Machine-readable ladder for docs/tools inventory."""
return {
"playbook_version": PLAYBOOK_VERSION,
"parent_issues": ["#655", "#652", "#653"],
"enforcement_issue": "#669",
"ladder": [r.as_dict() for r in RECOVERY_LADDER],
"broad_restart_actions": [a.value for a in sorted(BROAD_RESTART_ACTIONS, key=lambda x: x.value)],
"insufficient_outcomes": sorted(INSUFFICIENT_OUTCOMES),
"symptom_map": {k: v.value for k, v in sorted(SYMPTOM_TO_FIRST_ACTION.items())},
}
+46 -2
View File
@@ -1,4 +1,4 @@
"""MCP restart coordinator and impact analysis (#658). """MCP restart coordinator and impact analysis (#658 / #669).
Before any sanctioned MCP restart, a central coordinator must evaluate the Before any sanctioned MCP restart, a central coordinator must evaluate the
live control-plane state — active sessions, leases/locks, in-flight issue/PR live control-plane state — active sessions, leases/locks, in-flight issue/PR
@@ -16,6 +16,9 @@ Design rules (mirrors the read-only posture of ``workflow_dashboard`` /
a mutative apply path is a later child gated by a drain proof (non-goal here). a mutative apply path is a later child gated by a drain proof (non-goal here).
* **Fail closed.** If the inventory is not explicitly complete, the verdict is * **Fail closed.** If the inventory is not explicitly complete, the verdict is
``unsafe`` / deny — an incomplete evaluation must never green-light a restart. ``unsafe`` / deny — an incomplete evaluation must never green-light a restart.
* **Narrow-first (#669).** Broad classes (rolling / full / host) require a
prior attempt log of insufficient narrower recoveries unless break-glass is
authorized. See :mod:`recovery_playbook`.
* **No secrets.** Session ids, pids, and profiles are operational metadata, not * **No secrets.** Session ids, pids, and profiles are operational metadata, not
credentials; nothing secret flows through this module. credentials; nothing secret flows through this module.
@@ -32,8 +35,9 @@ from enum import Enum
from typing import Any, Mapping, Sequence from typing import Any, Mapping, Sequence
import lease_lifecycle import lease_lifecycle
import recovery_playbook
COORDINATOR_VERSION = "1.1.0-issue-663" COORDINATOR_VERSION = "1.2.0-issue-669"
# Restart verdicts. Exactly the three the acceptance criteria name. # Restart verdicts. Exactly the three the acceptance criteria name.
VERDICT_SAFE = "safe" VERDICT_SAFE = "safe"
@@ -349,6 +353,10 @@ class RestartImpactReport:
counts: dict[str, int] counts: dict[str, int]
audit_record: dict[str, Any] audit_record: dict[str, Any]
incomplete_reasons: list[str] = field(default_factory=list) incomplete_reasons: list[str] = field(default_factory=list)
# #669 playbook escalation gate (attempt-log enforcement).
playbook_escalation: dict[str, Any] = field(default_factory=dict)
attempt_log_satisfied: bool = True
break_glass: bool = False
def as_dict(self) -> dict[str, Any]: def as_dict(self) -> dict[str, Any]:
return { return {
@@ -382,6 +390,9 @@ class RestartImpactReport:
"prior_recovery_attempts": list(self.prior_recovery_attempts), "prior_recovery_attempts": list(self.prior_recovery_attempts),
"counts": dict(self.counts), "counts": dict(self.counts),
"audit_record": dict(self.audit_record), "audit_record": dict(self.audit_record),
"playbook_escalation": dict(self.playbook_escalation),
"attempt_log_satisfied": self.attempt_log_satisfied,
"break_glass": self.break_glass,
} }
@@ -496,6 +507,7 @@ def evaluate_restart_impact(
target_session_id: str | None = None, target_session_id: str | None = None,
target_role: str | None = None, target_role: str | None = None,
target_connector: str | None = None, target_connector: str | None = None,
break_glass: bool = False,
) -> RestartImpactReport: ) -> RestartImpactReport:
"""Evaluate a proposed MCP restart and return an impact preview. """Evaluate a proposed MCP restart and return an impact preview.
@@ -584,6 +596,30 @@ def evaluate_restart_impact(
dict(a) for a in (inventory.get("prior_recovery_attempts") or []) dict(a) for a in (inventory.get("prior_recovery_attempts") or [])
] ]
# #669: broad restarts require a prior narrow-attempt log unless break-glass.
playbook_escalation: dict[str, Any] = {}
attempt_log_satisfied = True
if policy_enforced and resolved_class is not None:
try:
escalation = recovery_playbook.assess_escalation(
resolved_class.value,
prior_recovery_attempts=prior_recovery_attempts,
break_glass=bool(break_glass),
)
playbook_escalation = escalation.as_dict()
attempt_log_satisfied = bool(escalation.allowed)
if not attempt_log_satisfied:
authorization_reasons.extend(list(escalation.reasons))
except ValueError as exc:
# Unknown mapping should never happen for enum values; fail closed.
attempt_log_satisfied = False
playbook_escalation = {
"allowed": False,
"reasons": [str(exc)],
"playbook_version": recovery_playbook.PLAYBOOK_VERSION,
}
authorization_reasons.append(str(exc))
session_impacts = [ session_impacts = [
_classify_session( _classify_session(
s, s,
@@ -682,6 +718,7 @@ def evaluate_restart_impact(
and role_authorized and role_authorized
and approval_satisfied and approval_satisfied
and target_complete and target_complete
and attempt_log_satisfied
) )
if policy_enforced and not authorization_ok: if policy_enforced and not authorization_ok:
@@ -745,6 +782,7 @@ def evaluate_restart_impact(
"affected_issues": len(affected_issues), "affected_issues": len(affected_issues),
"affected_prs": len(affected_prs), "affected_prs": len(affected_prs),
"prior_recovery_attempts": len(prior_recovery_attempts), "prior_recovery_attempts": len(prior_recovery_attempts),
"attempt_log_satisfied": attempt_log_satisfied,
} }
audit_record = { audit_record = {
@@ -765,6 +803,9 @@ def evaluate_restart_impact(
"allow_restart": allow_restart, "allow_restart": allow_restart,
"blast_radius": blast_radius, "blast_radius": blast_radius,
"counts": counts, "counts": counts,
"attempt_log_satisfied": attempt_log_satisfied,
"break_glass": bool(break_glass),
"playbook_version": recovery_playbook.PLAYBOOK_VERSION,
} }
return RestartImpactReport( return RestartImpactReport(
@@ -804,4 +845,7 @@ def evaluate_restart_impact(
counts=counts, counts=counts,
audit_record=audit_record, audit_record=audit_record,
incomplete_reasons=incomplete_reasons, incomplete_reasons=incomplete_reasons,
playbook_escalation=playbook_escalation,
attempt_log_satisfied=attempt_log_satisfied,
break_glass=bool(break_glass),
) )
@@ -35,6 +35,17 @@ BREAK_GLASS_ENV = "GITEA_BREAKGLASS_RESTART_AUTHORIZATION"
QUIET_SESSIONS: list[dict] = [] QUIET_SESSIONS: list[dict] = []
QUIET_LEASES: list[dict] = [] QUIET_LEASES: list[dict] = []
# #669: broad restarts need a prior narrow-attempt log (unless break-glass).
PRIOR_NARROW_ATTEMPTS_JSON = json.dumps(
[
{
"action": "client_reconnect",
"outcome": "insufficient",
"reason": "still flapping after reconnect",
}
]
)
class _FakeDB: class _FakeDB:
"""Minimal control-plane DB stand-in for the restart inventory.""" """Minimal control-plane DB stand-in for the restart inventory."""
@@ -128,6 +139,7 @@ class TestConjunction(_RestartToolHarness):
preview = self._call( preview = self._call(
role="operator", role="operator",
restart_class="full_mcp_restart", restart_class="full_mcp_restart",
prior_recovery_attempts_json=PRIOR_NARROW_ATTEMPTS_JSON,
env={CONTROLLER_APPROVAL_ENV: "operator-approved"}, env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
) )
self.assertTrue(preview["allow_restart"], self.assertTrue(preview["allow_restart"],
@@ -136,6 +148,7 @@ class TestConjunction(_RestartToolHarness):
result = self._call( result = self._call(
role="operator", role="operator",
restart_class="full_mcp_restart", restart_class="full_mcp_restart",
prior_recovery_attempts_json=PRIOR_NARROW_ATTEMPTS_JSON,
dry_run=False, dry_run=False,
drain_proof_json=self._clean_proof_for(preview), drain_proof_json=self._clean_proof_for(preview),
env={CONTROLLER_APPROVAL_ENV: "operator-approved"}, env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
@@ -342,6 +355,7 @@ class TestExistingPathsStillWork(_RestartToolHarness):
result = self._call( result = self._call(
role="operator", role="operator",
restart_class="full_mcp_restart", restart_class="full_mcp_restart",
prior_recovery_attempts_json=PRIOR_NARROW_ATTEMPTS_JSON,
dry_run=False, dry_run=False,
env={CONTROLLER_APPROVAL_ENV: "operator-approved"}, env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
) )
+217
View File
@@ -0,0 +1,217 @@
"""Unit tests for the scoped recovery playbook (#669)."""
from __future__ import annotations
import recovery_playbook as rp
import restart_coordinator as rc
def test_ladder_covers_eleven_ordered_rungs():
ranks = [r.rank for r in rp.RECOVERY_LADDER]
assert ranks == list(range(len(rp.RECOVERY_LADDER)))
assert len(rp.RECOVERY_LADDER) == 11
assert rp.RECOVERY_LADDER[0].action is rp.RecoveryAction.CLIENT_RECONNECT
assert rp.RECOVERY_LADDER[-1].action is rp.RecoveryAction.HOST_RESTART
def test_ladder_document_links_parent_issues():
doc = rp.ladder_document()
assert "#655" in doc["parent_issues"]
assert "#652" in doc["parent_issues"]
assert "#653" in doc["parent_issues"]
assert doc["enforcement_issue"] == "#669"
assert "full_mcp_restart" in doc["broad_restart_actions"]
def test_recommend_transport_eof_starts_at_client_reconnect():
plan = rp.recommend_actions(symptoms=["transport_eof"])
assert plan["recommended_actions"][0]["action"] == "client_reconnect"
assert plan["recommended_actions"][0]["issue_links"]
def test_recommend_skips_successful_prior_attempts():
attempts = [
rp.build_attempt_record(
"client_reconnect", outcome="success", reason="reconnected"
)
]
plan = rp.recommend_actions(
symptoms=["transport_eof"], prior_recovery_attempts=attempts
)
actions = [a["action"] for a in plan["recommended_actions"]]
assert "client_reconnect" not in actions
assert actions[0] == "capability_refresh"
def test_escalation_denied_without_attempt_log():
result = rp.assess_escalation("full_mcp_restart", prior_recovery_attempts=[])
assert result.allowed is False
assert result.require_attempt_log is True
assert any("#669" in r for r in result.reasons)
assert result.recommended_next # soft recommendations still provided
def test_escalation_allowed_after_insufficient_narrower():
attempts = [
rp.build_attempt_record(
"client_reconnect",
outcome="insufficient",
reason="still flapping",
),
rp.build_attempt_record(
"session_reconnect",
outcome="failed",
reason="namespace still dead",
),
]
result = rp.assess_escalation(
"full_mcp_restart", prior_recovery_attempts=attempts
)
assert result.allowed is True
assert len(result.qualifying_attempts) == 2
def test_escalation_break_glass_bypasses_attempt_log():
result = rp.assess_escalation(
"host_restart", prior_recovery_attempts=[], break_glass=True
)
assert result.allowed is True
assert result.break_glass is True
def test_narrow_action_does_not_require_attempt_log():
result = rp.assess_escalation(
"client_reconnect", prior_recovery_attempts=[]
)
assert result.allowed is True
assert result.require_attempt_log is False
def test_same_rank_attempt_does_not_qualify_for_escalation():
attempts = [
rp.build_attempt_record(
"full_mcp_restart", outcome="failed", reason="already failed full"
)
]
result = rp.assess_escalation(
"full_mcp_restart", prior_recovery_attempts=attempts
)
assert result.allowed is False
def test_success_outcome_does_not_qualify_for_escalation():
attempts = [
rp.build_attempt_record(
"client_reconnect", outcome="success", reason="fixed"
)
]
result = rp.assess_escalation(
"full_mcp_restart", prior_recovery_attempts=attempts
)
assert result.allowed is False
def test_recovery_metrics_fraction_avoided():
attempts = [
rp.build_attempt_record("client_reconnect", outcome="success"),
rp.build_attempt_record("session_reconnect", outcome="success"),
rp.build_attempt_record("full_mcp_restart", outcome="success"),
]
metrics = rp.recovery_metrics(attempts)
assert metrics["successes_total"] == 3
assert metrics["successes_avoided_full_restart"] == 2
assert metrics["successes_full_or_host_restart"] == 1
assert abs(metrics["fraction_avoided_full_restart"] - (2 / 3)) < 1e-9
def test_coordinator_denies_full_restart_without_attempt_log():
inv = {
"inventory_complete": True,
"sessions": [],
"leases": [],
"prior_recovery_attempts": [],
}
report = rc.evaluate_restart_impact(
inv,
restart_class=rc.RestartClass.FULL_MCP_RESTART,
requester_role="controller",
requester_permissions=rc.permissions_for_role("controller"),
controller_approved=True,
operator_authorized=True,
)
assert report.allow_restart is False
assert report.attempt_log_satisfied is False
assert report.verdict == rc.VERDICT_UNSAFE
blob = " ".join(report.reasons + report.authorization_reasons)
assert "#669" in blob or "attempt log" in blob
def test_coordinator_allows_full_restart_with_attempt_log():
inv = {
"inventory_complete": True,
"sessions": [],
"leases": [],
"prior_recovery_attempts": [
{
"action": "client_reconnect",
"outcome": "insufficient",
"reason": "still broken",
}
],
}
report = rc.evaluate_restart_impact(
inv,
restart_class=rc.RestartClass.FULL_MCP_RESTART,
requester_role="controller",
requester_permissions=rc.permissions_for_role("controller"),
controller_approved=True,
operator_authorized=True,
)
assert report.attempt_log_satisfied is True
assert report.allow_restart is True
assert report.verdict == rc.VERDICT_SAFE
def test_coordinator_break_glass_allows_without_log():
inv = {
"inventory_complete": True,
"sessions": [],
"leases": [],
"prior_recovery_attempts": [],
}
report = rc.evaluate_restart_impact(
inv,
restart_class=rc.RestartClass.FULL_MCP_RESTART,
requester_role="controller",
requester_permissions=rc.permissions_for_role("controller"),
controller_approved=True,
operator_authorized=True,
break_glass=True,
)
assert report.break_glass is True
assert report.attempt_log_satisfied is True
assert report.allow_restart is True
def test_coordinator_client_reconnect_unaffected():
inv = {
"inventory_complete": True,
"sessions": [],
"leases": [],
"prior_recovery_attempts": [],
}
report = rc.evaluate_restart_impact(
inv,
restart_class=rc.RestartClass.CLIENT_RECONNECT,
requester_role="author",
requester_permissions=rc.permissions_for_role("author"),
)
assert report.attempt_log_satisfied is True
assert report.allow_restart is True
def test_restart_class_alias_accepted():
assert (
rp.resolve_action("full_mcp_restart")
is rp.RecoveryAction.FULL_MCP_RESTART
)
-404
View File
@@ -1,404 +0,0 @@
"""Tests for AI-provider connections and evidence-backed insights (#650).
Covers acceptance criteria:
1. Provider connection status is redacted and accurate (declared registry only).
2. At least three insight types with evidence citations.
3. Insights never claim actions completed without proof.
4. Evidence requirement is enforced (no evidence-less insights).
5. Interpretation limits appear in docs-facing payloads.
"""
from __future__ import annotations
import json
import sys
import unittest
from dataclasses import dataclass
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from tests.webui_testclient import TestClient
from webui.app import create_app
from webui.insights_loader import (
CONFIDENCE_HIGH,
CONNECTION_DECLARED_AVAILABLE,
CONNECTION_DECLARED_UNAVAILABLE,
INSIGHT_BLOCKED_QUEUE,
INSIGHT_CONTROLLER_ATTENTION,
INSIGHT_PROVIDER_WITHOUT_WORKERS,
INSIGHT_STALE_RUNTIME,
ProviderConnection,
build_provider_connection,
generate_insights,
insight_blocked_queue,
insight_controller_attention,
insight_providers_without_workers,
insight_stale_runtime,
load_insights_snapshot,
load_provider_snapshot,
snapshot_insights_to_dict,
snapshot_providers_to_dict,
)
from webui.insights_views import render_insights_page, render_providers_page
from webui.nav import STUB_PAGES, nav_hrefs
from webui.worker_registry import ProviderRecord, WorkerRecord, ScheduleSpec, SchedulerSpec
def _provider(
provider_id: str = "claude",
*,
available: bool = True,
models: tuple[str, ...] = ("claude-opus-4-8",),
notes: str = "",
) -> ProviderRecord:
return ProviderRecord(
id=provider_id,
display_name=provider_id.title(),
vendor="TestVendor",
executable=provider_id,
available=available,
models=models,
notes=notes,
)
def _worker(
worker_id: str = "claude-author",
*,
provider: str = "claude",
enabled: bool = True,
) -> WorkerRecord:
return WorkerRecord(
id=worker_id,
display_name=worker_id,
provider=provider,
model="m1",
project="gitea-tools",
role="author",
namespace="gitea-author",
profile="prgs-author",
workflow="skills/llm-project-workflow/workflows/work-issue.md",
schedule=ScheduleSpec(kind="manual", seconds=None, expression=None),
timeout_seconds=3600,
enabled=enabled,
scheduler=SchedulerSpec(kind="manual", label=None),
notes="",
)
@dataclass(frozen=True)
class _TrafficItem:
kind: str
number: int
title: str = ""
traffic_state: str = "blocked"
expected_role: str = "author"
safe_for_roles: tuple[str, ...] = ()
badges: tuple[str, ...] = ()
block_reason: str | None = "dependency"
@dataclass(frozen=True)
class _Traffic:
blocked: tuple = ()
needs_controller: tuple = ()
inventory_complete: bool = True
fetch_error: str | None = None
@dataclass(frozen=True)
class _Stale:
daemon_head: str | None
checkout_head: str | None
remote_head: str | None
stale: bool
determinable: bool
mutation_safe: bool
@dataclass(frozen=True)
class _Health:
stale_runtime: _Stale | None
class TestProviderConnections(unittest.TestCase):
def test_available_provider_status(self):
conn = build_provider_connection(_provider(available=True), (_worker(),))
self.assertEqual(conn.connection_status, CONNECTION_DECLARED_AVAILABLE)
self.assertTrue(conn.available_declared)
self.assertEqual(conn.worker_count, 1)
self.assertEqual(conn.enabled_worker_count, 1)
self.assertFalse(conn.to_dict()["secrets_exposed"])
def test_unavailable_provider_status(self):
conn = build_provider_connection(_provider(available=False), ())
self.assertEqual(conn.connection_status, CONNECTION_DECLARED_UNAVAILABLE)
self.assertEqual(conn.worker_count, 0)
def test_notes_are_redacted(self):
conn = build_provider_connection(
_provider(notes="token=ghp_thisisnotarealsecretvalue0001"),
(),
)
self.assertNotIn("ghp_thisisnotarealsecretvalue0001", conn.notes)
self.assertNotIn(
"ghp_thisisnotarealsecretvalue0001",
json.dumps(conn.to_dict()),
)
def test_load_provider_snapshot_from_injected_registry(self):
from webui.worker_registry import WorkerRegistry
from pathlib import Path
registry = WorkerRegistry(
version=1,
revision=3,
updated_at="2026-07-25T00:00:00Z",
providers=(_provider("claude"), _provider("grok", available=False)),
workers=(_worker(),),
source_path=Path("/tmp/workers.registry.json"),
)
snapshot = load_provider_snapshot(registry=registry)
self.assertTrue(snapshot.ok)
self.assertEqual(snapshot.registry_revision, 3)
ids = {p.provider_id for p in snapshot.providers}
self.assertEqual(ids, {"claude", "grok"})
def test_registry_failure_is_fail_closed(self):
def _boom():
raise RuntimeError("disk gone")
snapshot = load_provider_snapshot(registry_loader=_boom)
self.assertFalse(snapshot.ok)
self.assertIn("unavailable", snapshot.fetch_error or "")
self.assertEqual(snapshot.providers, ())
class TestInsightGenerators(unittest.TestCase):
def test_blocked_queue_requires_evidence(self):
traffic = _Traffic(
blocked=(
_TrafficItem(kind="issue", number=647, block_reason="depends #646"),
_TrafficItem(kind="pr", number=902, block_reason="conflict"),
)
)
insight = insight_blocked_queue(traffic)
self.assertIsNotNone(insight)
self.assertEqual(insight.kind, INSIGHT_BLOCKED_QUEUE)
self.assertGreaterEqual(len(insight.evidence), 2)
self.assertTrue(insight.advisory_only)
self.assertFalse(insight.claims_action_completed)
refs = {e.ref for e in insight.evidence}
self.assertIn("#647", refs)
self.assertIn("#902", refs)
def test_empty_blocked_queue_yields_no_insight(self):
self.assertIsNone(insight_blocked_queue(_Traffic()))
def test_controller_attention_insight(self):
traffic = _Traffic(
needs_controller=(_TrafficItem(kind="issue", number=100, traffic_state="needs_controller"),)
)
insight = insight_controller_attention(traffic)
self.assertEqual(insight.kind, INSIGHT_CONTROLLER_ATTENTION)
self.assertEqual(insight.evidence[0].ref, "#100")
def test_stale_runtime_insight(self):
health = _Health(
stale_runtime=_Stale(
daemon_head="aaa",
checkout_head="bbb",
remote_head="ccc",
stale=True,
determinable=True,
mutation_safe=False,
)
)
insight = insight_stale_runtime(health)
self.assertEqual(insight.kind, INSIGHT_STALE_RUNTIME)
self.assertIn("stale", insight.evidence[0].detail)
self.assertFalse(insight.claims_action_completed)
def test_mutation_safe_runtime_yields_no_insight(self):
health = _Health(
stale_runtime=_Stale(
daemon_head="aaa",
checkout_head="aaa",
remote_head="aaa",
stale=False,
determinable=True,
mutation_safe=True,
)
)
self.assertIsNone(insight_stale_runtime(health))
def test_provider_without_workers(self):
providers = (
build_provider_connection(_provider("claude"), (_worker(),)),
build_provider_connection(_provider("grok"), ()),
)
insight = insight_providers_without_workers(providers)
self.assertEqual(insight.kind, INSIGHT_PROVIDER_WITHOUT_WORKERS)
self.assertEqual(insight.evidence[0].ref, "grok")
def test_generate_insights_composes_three_kinds(self):
from webui.worker_registry import WorkerRegistry
from pathlib import Path
registry = WorkerRegistry(
version=1,
revision=1,
updated_at="2026-07-25T00:00:00Z",
providers=(_provider("lonely"),),
workers=(),
source_path=Path("/tmp/w.json"),
)
provider_snapshot = load_provider_snapshot(registry=registry)
traffic = _Traffic(
blocked=(_TrafficItem(kind="issue", number=1),),
needs_controller=(_TrafficItem(kind="issue", number=2),),
)
health = _Health(
stale_runtime=_Stale("a", "b", "c", True, True, False)
)
insights, used, unavailable = generate_insights(
traffic=traffic,
health=health,
provider_snapshot=provider_snapshot,
analytics=None,
)
kinds = {i.kind for i in insights}
self.assertIn(INSIGHT_BLOCKED_QUEUE, kinds)
self.assertIn(INSIGHT_CONTROLLER_ATTENTION, kinds)
self.assertIn(INSIGHT_STALE_RUNTIME, kinds)
self.assertIn(INSIGHT_PROVIDER_WITHOUT_WORKERS, kinds)
self.assertGreaterEqual(len(kinds), 3)
self.assertIn("traffic", used)
self.assertIn("system_health", used)
self.assertIn("providers", used)
self.assertTrue(any(u["source"] == "analytics" for u in unavailable))
for insight in insights:
self.assertTrue(insight.advisory_only)
self.assertFalse(insight.claims_action_completed)
self.assertGreaterEqual(len(insight.evidence), 1)
def test_missing_source_is_reported_not_as_healthy_empty(self):
insights, used, unavailable = generate_insights(
traffic=None,
health=None,
provider_snapshot=None,
analytics=None,
)
self.assertEqual(insights, ())
self.assertEqual(used, ())
self.assertEqual(len(unavailable), 4)
class TestViewsAndRoutes(unittest.TestCase):
def test_providers_page_renders_connections(self):
from webui.worker_registry import WorkerRegistry
from pathlib import Path
registry = WorkerRegistry(
version=1,
revision=1,
updated_at="2026-07-25T00:00:00Z",
providers=(_provider("claude"),),
workers=(_worker(),),
source_path=Path("/tmp/w.json"),
)
snapshot = load_provider_snapshot(registry=registry)
html = render_providers_page(snapshot)
self.assertIn("AI provider connections", html)
self.assertIn("claude", html)
self.assertIn("declared_available", html)
self.assertIn("Interpretation limits", html)
self.assertNotIn("api_key", html.lower())
def test_failed_provider_snapshot_renders_no_table(self):
snapshot = load_provider_snapshot(registry_loader=lambda: (_ for _ in ()).throw(RuntimeError("x")))
html = render_providers_page(snapshot)
self.assertIn("unavailable", html.lower())
self.assertNotIn("<tbody><tr><td><code>", html)
def test_insights_page_lists_evidence(self):
traffic = _Traffic(blocked=(_TrafficItem(kind="issue", number=42),))
snapshot = load_insights_snapshot(
traffic=traffic,
health=_Health(None),
provider_snapshot=load_provider_snapshot(
registry_loader=lambda: (_ for _ in ()).throw(RuntimeError("skip"))
),
analytics=None,
load_live=False,
)
html = render_insights_page(snapshot)
self.assertIn("#42", html)
self.assertIn("advisory only", html.lower())
self.assertIn("Evidence", html)
def test_nav_exposes_live_insights_and_providers(self):
self.assertIn("/insights", nav_hrefs())
self.assertIn("/providers", nav_hrefs())
self.assertNotIn("/insights", STUB_PAGES)
def test_routes_are_read_only_and_export_json(self):
client = TestClient(create_app())
# Use live registry from package data — should be ok.
with mock.patch(
"webui.app.load_provider_snapshot",
return_value=load_provider_snapshot(
registry=__import__(
"webui.worker_registry", fromlist=["load_registry"]
).load_registry()
),
):
response = client.get("/providers")
self.assertEqual(response.status_code, 200)
self.assertIn("provider", response.text.lower())
api = client.get("/api/v1/providers")
self.assertEqual(api.status_code, 200)
payload = api.json()
self.assertTrue(payload["ok"])
self.assertIn("interpretation_limits", payload)
self.assertTrue(all(not p.get("secrets_exposed") for p in payload["providers"]))
with mock.patch(
"webui.app.load_insights_snapshot",
return_value=load_insights_snapshot(
traffic=_Traffic(blocked=(_TrafficItem(kind="issue", number=7),)),
health=_Health(None),
provider_snapshot=load_provider_snapshot(
registry_loader=lambda: (_ for _ in ()).throw(RuntimeError("x"))
),
analytics=None,
load_live=False,
),
):
page = client.get("/insights")
self.assertEqual(page.status_code, 200)
self.assertIn("#7", page.text)
api = client.get("/api/v1/insights")
self.assertEqual(api.status_code, 200)
body = api.json()
self.assertTrue(body["ok"])
self.assertTrue(all(i["advisory_only"] for i in body["insights"]))
self.assertTrue(all(not i["claims_action_completed"] for i in body["insights"]))
self.assertTrue(all(i["evidence"] for i in body["insights"]))
for path in ("/providers", "/api/v1/providers", "/insights", "/api/v1/insights"):
with self.subTest(path=path):
self.assertEqual(client.post(path).status_code, 405)
def test_home_nav_links_providers_and_insights(self):
home = TestClient(create_app()).get("/").text
self.assertIn('href="/providers"', home)
self.assertIn('href="/insights"', home)
if __name__ == "__main__":
unittest.main()
+2 -39
View File
@@ -53,13 +53,6 @@ from webui.session_loader import (
snapshot_to_dict as session_view_snapshot_to_dict, snapshot_to_dict as session_view_snapshot_to_dict,
) )
from webui.session_views import render_sessions_page from webui.session_views import render_sessions_page
from webui.insights_loader import (
load_insights_snapshot,
load_provider_snapshot,
snapshot_insights_to_dict,
snapshot_providers_to_dict,
)
from webui.insights_views import render_insights_page, render_providers_page
from webui.linkage_loader import ( from webui.linkage_loader import (
load_linkage_snapshot, load_linkage_snapshot,
snapshot_to_dict as linkage_snapshot_to_dict, snapshot_to_dict as linkage_snapshot_to_dict,
@@ -354,34 +347,6 @@ async def api_sessions(_request: Request) -> JSONResponse:
return JSONResponse(session_view_snapshot_to_dict(load_session_view_snapshot())) return JSONResponse(session_view_snapshot_to_dict(load_session_view_snapshot()))
async def providers(_request: Request) -> HTMLResponse:
"""AI-provider connection status (#650) — declared registry only, no secrets."""
return HTMLResponse(render_providers_page(load_provider_snapshot()))
async def api_v1_providers(_request: Request) -> JSONResponse:
"""JSON export of declared AI-provider connections (#650)."""
snapshot = load_provider_snapshot()
return JSONResponse(
snapshot_providers_to_dict(snapshot),
status_code=200 if snapshot.ok else 502,
)
async def insights(_request: Request) -> HTMLResponse:
"""Evidence-backed operational insights (#650) — advisory only."""
return HTMLResponse(render_insights_page(load_insights_snapshot()))
async def api_v1_insights(_request: Request) -> JSONResponse:
"""JSON export of evidence-backed insights (#650)."""
snapshot = load_insights_snapshot()
return JSONResponse(
snapshot_insights_to_dict(snapshot),
status_code=200 if snapshot.ok else 502,
)
def _linkage_snapshot(request: Request): def _linkage_snapshot(request: Request):
"""Load one linkage snapshot from the request's scope and focus parameters.""" """Load one linkage snapshot from the request's scope and focus parameters."""
return load_linkage_snapshot( return load_linkage_snapshot(
@@ -415,6 +380,8 @@ async def api_v1_gitea_linkage(request: Request) -> JSONResponse:
linkage_snapshot_to_dict(snapshot), linkage_snapshot_to_dict(snapshot),
status_code=200 if snapshot.ok else 502, status_code=200 if snapshot.ok else 502,
) )
async def _parse_audit_form(request: Request) -> tuple[str, str | None]: async def _parse_audit_form(request: Request) -> tuple[str, str | None]:
if request.method == "GET": if request.method == "GET":
return "", None return "", None
@@ -858,10 +825,6 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/api/sessions", api_sessions, methods=["GET"]), Route("/api/sessions", api_sessions, methods=["GET"]),
Route("/api/v1/sessions", api_sessions, methods=["GET"]), Route("/api/v1/sessions", api_sessions, methods=["GET"]),
Route("/api/v1/timeline", api_v1_timeline, methods=["GET"]), Route("/api/v1/timeline", api_v1_timeline, methods=["GET"]),
Route("/providers", providers, methods=["GET"]),
Route("/api/v1/providers", api_v1_providers, methods=["GET"]),
Route("/insights", insights, methods=["GET"]),
Route("/api/v1/insights", api_v1_insights, methods=["GET"]),
Route("/gitea", gitea_linkage, methods=["GET"]), Route("/gitea", gitea_linkage, methods=["GET"]),
Route("/api/v1/gitea/linkage", api_v1_gitea_linkage, methods=["GET"]), Route("/api/v1/gitea/linkage", api_v1_gitea_linkage, methods=["GET"]),
Route("/analytics", analytics, methods=["GET"]), Route("/analytics", analytics, methods=["GET"]),
-713
View File
@@ -1,713 +0,0 @@
"""AI-provider connections and evidence-backed operational insights (#650, Phase 4).
Operators need two related, **advisory** surfaces:
1. **Provider connection status** — which AI runtimes are *declared* in the
worker registry (#798), without ever exposing API keys or inventing a live
probe that this process cannot perform.
2. **Evidence-backed insights** — short cards derived only from durable
console evidence (traffic, system health, analytics, the same registry).
Every insight carries explicit evidence refs (issue/PR/provider/event ids).
Insights never claim that a workflow action completed without proof, and
they never mutate anything.
Design rules matching the rest of the console:
- **Read-only.** No endpoint registered here mutates Gitea, the control plane,
or the registry.
- **Advisory only.** Insights carry ``advisory_only=True`` and never emit an
"action completed" claim. The allocator, review, and merge paths remain the
only authorities for work selection and terminal state.
- **Qualified absence.** When a source could not run, the insight list says so
rather than inventing an empty-and-healthy fleet or zero blocked items.
- **Redaction.** Free-text titles, reasons, and notes pass through
``webui.console_redaction`` before they leave this module.
- **No secrets.** Provider records are taken from the credential-free worker
registry. Keys never appear in this surface.
Non-goals (from the issue): free-form chatbot that overrides gates, secret
provider keys in the UI, auto-merge or auto-close from insights.
"""
from __future__ import annotations
import os
from dataclasses import dataclass
from typing import Any, Callable, Sequence
from webui import console_redaction
from webui.worker_registry import (
ProviderRecord,
WorkerRegistry,
WorkerRecord,
load_registry as load_worker_registry,
workers_for_provider,
)
INSIGHTS_SCHEMA_VERSION = 1
# Provider connection vocabulary. Declared availability is not a live probe —
# the worker registry owns the declaration, and adapters (#800) own live checks.
CONNECTION_DECLARED_AVAILABLE = "declared_available"
CONNECTION_DECLARED_UNAVAILABLE = "declared_unavailable"
CONNECTION_REGISTRY_UNAVAILABLE = "registry_unavailable"
# Insight kinds. Each generator is a pure function over one evidence source.
INSIGHT_BLOCKED_QUEUE = "blocked_queue_pressure"
INSIGHT_CONTROLLER_ATTENTION = "controller_attention"
INSIGHT_STALE_RUNTIME = "stale_runtime_risk"
INSIGHT_PROVIDER_WITHOUT_WORKERS = "provider_without_workers"
INSIGHT_ANALYTICS_FAILURE_RATE = "analytics_failure_pressure"
SEVERITY_INFO = "info"
SEVERITY_WARN = "warn"
SEVERITY_CRITICAL = "critical"
SEVERITY_UNPROVEN = "unproven"
CONFIDENCE_HIGH = "high"
CONFIDENCE_MEDIUM = "medium"
CONFIDENCE_LOW = "low"
CONFIDENCE_UNPROVEN = "unproven"
def _redact(value: Any) -> Any:
if value is None:
return None
return console_redaction.redact_text(str(value))
def _offline_test_mode() -> bool:
return (os.environ.get("WEBUI_TEST_OFFLINE") or "").strip().lower() in {
"1",
"true",
"yes",
"on",
}
# --- Provider connection status ------------------------------------------------
@dataclass(frozen=True)
class ProviderConnection:
"""One AI provider's declared connection status (no secrets, no live probe)."""
provider_id: str
display_name: str
vendor: str
executable: str
connection_status: str
available_declared: bool
models: tuple[str, ...]
worker_count: int
enabled_worker_count: int
notes: str
#: Explicit statement of what was *not* proven (live process health, etc.).
probe_limit: str
def to_dict(self) -> dict[str, Any]:
return {
"provider_id": self.provider_id,
"display_name": self.display_name,
"vendor": self.vendor,
"executable": self.executable,
"connection_status": self.connection_status,
"available_declared": self.available_declared,
"models": list(self.models),
"worker_count": self.worker_count,
"enabled_worker_count": self.enabled_worker_count,
"notes": self.notes,
"probe_limit": self.probe_limit,
# Always true for this surface: keys are never loaded.
"secrets_exposed": False,
}
@dataclass(frozen=True)
class ProviderSnapshot:
ok: bool
providers: tuple[ProviderConnection, ...] = ()
registry_revision: int | None = None
registry_path: str | None = None
fetch_error: str | None = None
schema_version: int = INSIGHTS_SCHEMA_VERSION
def to_dict(self) -> dict[str, Any]:
return {
"ok": self.ok,
"schema_version": self.schema_version,
"registry_revision": self.registry_revision,
"registry_path": self.registry_path,
"fetch_error": self.fetch_error,
"providers": [p.to_dict() for p in self.providers],
"interpretation_limits": [
"connection_status reflects the worker registry declaration only",
"no API keys or credential material are loaded or rendered",
"live executable health is not probed on this surface (#800 owns that)",
],
}
_PROBE_LIMIT = (
"Declared status only. This console does not probe the provider executable "
"or call vendor APIs; live health belongs to the provider adapter framework."
)
def connection_status_for(provider: ProviderRecord) -> str:
return (
CONNECTION_DECLARED_AVAILABLE
if provider.available
else CONNECTION_DECLARED_UNAVAILABLE
)
def build_provider_connection(
provider: ProviderRecord,
workers: Sequence[WorkerRecord],
) -> ProviderConnection:
enabled = sum(1 for worker in workers if worker.enabled)
return ProviderConnection(
provider_id=provider.id,
display_name=str(_redact(provider.display_name) or provider.id),
vendor=str(_redact(provider.vendor) or ""),
executable=str(_redact(provider.executable) or ""),
connection_status=connection_status_for(provider),
available_declared=bool(provider.available),
models=tuple(str(_redact(m) or m) for m in provider.models),
worker_count=len(workers),
enabled_worker_count=enabled,
notes=str(_redact(provider.notes) or ""),
probe_limit=_PROBE_LIMIT,
)
def load_provider_snapshot(
*,
registry: WorkerRegistry | None = None,
registry_loader: Callable[[], WorkerRegistry] | None = None,
) -> ProviderSnapshot:
"""Load declared provider connections. Never raises for missing registry."""
if registry is None:
loader = registry_loader or load_worker_registry
try:
if _offline_test_mode() and registry_loader is None:
return ProviderSnapshot(
ok=False,
fetch_error=(
"provider registry not loaded in offline test mode "
"(inject a registry for unit tests)"
),
)
registry = loader()
except Exception as exc: # fail soft — operator-visible reason
return ProviderSnapshot(
ok=False,
fetch_error=str(_redact(f"worker registry unavailable: {exc}")),
)
connections = tuple(
build_provider_connection(provider, workers_for_provider(registry, provider.id))
for provider in registry.providers
)
return ProviderSnapshot(
ok=True,
providers=connections,
registry_revision=registry.revision,
registry_path=str(registry.source_path),
)
# --- Evidence-backed insights --------------------------------------------------
@dataclass(frozen=True)
class EvidenceRef:
"""One durable reference an insight is allowed to cite."""
kind: str # issue | pr | provider | health | analytics | traffic
ref: str
detail: str
def to_dict(self) -> dict[str, Any]:
return {
"kind": self.kind,
"ref": self.ref,
"detail": str(_redact(self.detail) or ""),
}
@dataclass(frozen=True)
class Insight:
"""One advisory finding. Never a claim that an action completed."""
insight_id: str
kind: str
severity: str
confidence: str
title: str
summary: str
evidence: tuple[EvidenceRef, ...]
advisory_only: bool = True
claims_action_completed: bool = False
def to_dict(self) -> dict[str, Any]:
return {
"insight_id": self.insight_id,
"kind": self.kind,
"severity": self.severity,
"confidence": self.confidence,
"title": str(_redact(self.title) or ""),
"summary": str(_redact(self.summary) or ""),
"evidence": [item.to_dict() for item in self.evidence],
"advisory_only": self.advisory_only,
"claims_action_completed": self.claims_action_completed,
}
@dataclass(frozen=True)
class InsightsSnapshot:
ok: bool
insights: tuple[Insight, ...] = ()
sources_used: tuple[str, ...] = ()
sources_unavailable: tuple[dict[str, str], ...] = ()
fetch_error: str | None = None
schema_version: int = INSIGHTS_SCHEMA_VERSION
def to_dict(self) -> dict[str, Any]:
return {
"ok": self.ok,
"schema_version": self.schema_version,
"insights": [insight.to_dict() for insight in self.insights],
"sources_used": list(self.sources_used),
"sources_unavailable": list(self.sources_unavailable),
"fetch_error": self.fetch_error,
"interpretation_limits": [
"insights are advisory only and never authorize merge, review, or close",
"an insight without evidence refs is refused rather than emitted",
"a missing source is listed under sources_unavailable, not as an empty success",
"insights never claim a workflow action completed",
],
}
def _require_evidence(evidence: Sequence[EvidenceRef]) -> tuple[EvidenceRef, ...]:
"""Fail closed: an insight with no evidence must not be emitted."""
items = tuple(evidence)
if not items:
raise ValueError("insight requires at least one evidence ref")
return items
def insight_blocked_queue(traffic: Any) -> Insight | None:
"""Traffic blocked bucket pressure with per-item evidence."""
blocked = tuple(getattr(traffic, "blocked", ()) or ())
if not blocked:
return None
evidence = []
for item in blocked[:20]:
kind = str(getattr(item, "kind", "issue") or "issue")
number = int(getattr(item, "number", 0) or 0)
if number <= 0:
continue
reason = getattr(item, "block_reason", None) or "blocked"
evidence.append(
EvidenceRef(
kind=kind,
ref=f"#{number}",
detail=f"traffic_state=blocked; reason={reason}",
)
)
if not evidence:
return None
count = len(blocked)
severity = SEVERITY_CRITICAL if count >= 10 else SEVERITY_WARN
return Insight(
insight_id=f"{INSIGHT_BLOCKED_QUEUE}:{count}",
kind=INSIGHT_BLOCKED_QUEUE,
severity=severity,
confidence=(
CONFIDENCE_HIGH
if getattr(traffic, "inventory_complete", False)
else CONFIDENCE_MEDIUM
),
title=f"{count} blocked work item(s) in traffic control",
summary=(
f"Traffic control reports {count} blocked item(s). "
"This is an observation of the loaded window, not a claim that "
"any remediation ran."
),
evidence=_require_evidence(evidence),
)
def insight_controller_attention(traffic: Any) -> Insight | None:
needs = tuple(getattr(traffic, "needs_controller", ()) or ())
if not needs:
return None
evidence = []
for item in needs[:20]:
kind = str(getattr(item, "kind", "issue") or "issue")
number = int(getattr(item, "number", 0) or 0)
if number <= 0:
continue
evidence.append(
EvidenceRef(
kind=kind,
ref=f"#{number}",
detail="traffic_state=needs_controller",
)
)
if not evidence:
return None
count = len(needs)
return Insight(
insight_id=f"{INSIGHT_CONTROLLER_ATTENTION}:{count}",
kind=INSIGHT_CONTROLLER_ATTENTION,
severity=SEVERITY_WARN if count else SEVERITY_INFO,
confidence=(
CONFIDENCE_HIGH
if getattr(traffic, "inventory_complete", False)
else CONFIDENCE_MEDIUM
),
title=f"{count} item(s) need controller attention",
summary=(
f"Traffic control marks {count} item(s) as needs_controller. "
"Advisory only — the controller allocator remains the authority "
"for routing."
),
evidence=_require_evidence(evidence),
)
def insight_stale_runtime(health: Any) -> Insight | None:
stale = getattr(health, "stale_runtime", None)
if stale is None:
return None
mutation_safe = bool(getattr(stale, "mutation_safe", False))
is_stale = bool(getattr(stale, "stale", False))
determinable = bool(getattr(stale, "determinable", False))
if mutation_safe and not is_stale:
return None
daemon = getattr(stale, "daemon_head", None) or "unknown"
checkout = getattr(stale, "checkout_head", None) or "unknown"
remote = getattr(stale, "remote_head", None) or "unknown"
if not determinable:
severity = SEVERITY_UNPROVEN
confidence = CONFIDENCE_UNPROVEN
title = "Runtime parity is not determinable"
summary = (
"System health could not prove mutation_safe. This is not proof "
"that the runtime is stale — only that parity was unproven."
)
else:
severity = SEVERITY_CRITICAL if is_stale else SEVERITY_WARN
confidence = CONFIDENCE_HIGH
title = "Stale or mutation-unsafe runtime"
summary = (
"System health reports a runtime that is not mutation_safe. "
"No restart or recovery is claimed by this insight."
)
return Insight(
insight_id=f"{INSIGHT_STALE_RUNTIME}:{daemon}:{checkout}",
kind=INSIGHT_STALE_RUNTIME,
severity=severity,
confidence=confidence,
title=title,
summary=summary,
evidence=_require_evidence(
(
EvidenceRef(
kind="health",
ref="stale_runtime",
detail=(
f"stale={is_stale}; mutation_safe={mutation_safe}; "
f"determinable={determinable}; daemon={daemon}; "
f"checkout={checkout}; remote={remote}"
),
),
)
),
)
def insight_providers_without_workers(
providers: Sequence[ProviderConnection],
) -> Insight | None:
lonely = [
provider
for provider in providers
if provider.available_declared and provider.worker_count == 0
]
if not lonely:
return None
evidence = tuple(
EvidenceRef(
kind="provider",
ref=provider.provider_id,
detail=(
f"available_declared=true; worker_count=0; "
f"vendor={provider.vendor}"
),
)
for provider in lonely
)
return Insight(
insight_id=f"{INSIGHT_PROVIDER_WITHOUT_WORKERS}:{len(lonely)}",
kind=INSIGHT_PROVIDER_WITHOUT_WORKERS,
severity=SEVERITY_INFO,
confidence=CONFIDENCE_HIGH,
title=f"{len(lonely)} declared-available provider(s) have no workers",
summary=(
"The worker registry declares these providers available but no "
"worker instance names them. This is a configuration observation, "
"not a claim that a provider process is running or idle."
),
evidence=_require_evidence(evidence),
)
def insight_analytics_failures(analytics: Any) -> Insight | None:
"""Flag elevated non-ok stage status in analytics when events exist."""
if analytics is None or not getattr(analytics, "ok", False):
return None
events = tuple(getattr(analytics, "events", ()) or ())
if not events:
return None
failed = [
event
for event in events
if str(getattr(event, "status", "") or "").lower()
in {"error", "failed", "failure"}
]
if not failed:
return None
# Cap evidence so a large window stays readable.
evidence = []
for event in failed[:20]:
usage_id = getattr(event, "usage_id", None)
issue = getattr(event, "issue_number", None)
pr = getattr(event, "pr_number", None)
if pr is not None:
ref_kind, ref = "pr", f"#{int(pr)}"
elif issue is not None:
ref_kind, ref = "issue", f"#{int(issue)}"
else:
ref_kind, ref = "analytics", f"usage:{usage_id}"
evidence.append(
EvidenceRef(
kind=ref_kind,
ref=ref,
detail=(
f"status={getattr(event, 'status', '')}; "
f"stage={getattr(event, 'stage', '')}; "
f"model={getattr(event, 'model', '')}"
),
)
)
if not evidence:
return None
rate = len(failed) / max(len(events), 1)
return Insight(
insight_id=f"{INSIGHT_ANALYTICS_FAILURE_RATE}:{len(failed)}:{len(events)}",
kind=INSIGHT_ANALYTICS_FAILURE_RATE,
severity=SEVERITY_WARN if rate >= 0.1 else SEVERITY_INFO,
confidence=CONFIDENCE_MEDIUM,
title=f"{len(failed)} analytics event(s) reported failure status",
summary=(
f"{len(failed)} of {len(events)} loaded analytics events carry a "
"failure status. Advisory only — this is not a gate decision."
),
evidence=_require_evidence(evidence),
)
def generate_insights(
*,
traffic: Any | None = None,
health: Any | None = None,
provider_snapshot: ProviderSnapshot | None = None,
analytics: Any | None = None,
) -> tuple[tuple[Insight, ...], tuple[str, ...], tuple[dict[str, str], ...]]:
"""Pure multi-source insight generation. Never mutates inputs."""
insights: list[Insight] = []
used: list[str] = []
unavailable: list[dict[str, str]] = []
if traffic is None:
unavailable.append(
{"source": "traffic", "reason": "traffic snapshot not supplied"}
)
elif getattr(traffic, "fetch_error", None):
unavailable.append(
{
"source": "traffic",
"reason": str(_redact(traffic.fetch_error) or "traffic fetch failed"),
}
)
else:
used.append("traffic")
for builder in (insight_blocked_queue, insight_controller_attention):
try:
item = builder(traffic)
except ValueError:
continue
if item is not None:
insights.append(item)
if health is None:
unavailable.append(
{"source": "system_health", "reason": "system health snapshot not supplied"}
)
else:
used.append("system_health")
try:
item = insight_stale_runtime(health)
except ValueError:
item = None
if item is not None:
insights.append(item)
if provider_snapshot is None:
unavailable.append(
{"source": "providers", "reason": "provider snapshot not supplied"}
)
elif not provider_snapshot.ok:
unavailable.append(
{
"source": "providers",
"reason": str(
_redact(provider_snapshot.fetch_error)
or "provider registry unavailable"
),
}
)
else:
used.append("providers")
try:
item = insight_providers_without_workers(provider_snapshot.providers)
except ValueError:
item = None
if item is not None:
insights.append(item)
if analytics is None:
unavailable.append(
{"source": "analytics", "reason": "analytics snapshot not supplied"}
)
elif not getattr(analytics, "ok", False):
unavailable.append(
{
"source": "analytics",
"reason": str(
_redact(getattr(analytics, "fetch_error", None))
or "analytics snapshot not ok"
),
}
)
else:
used.append("analytics")
try:
item = insight_analytics_failures(analytics)
except ValueError:
item = None
if item is not None:
insights.append(item)
# Stable ordering: severity then kind.
_sev_rank = {
SEVERITY_CRITICAL: 0,
SEVERITY_WARN: 1,
SEVERITY_INFO: 2,
SEVERITY_UNPROVEN: 3,
}
insights.sort(key=lambda i: (_sev_rank.get(i.severity, 9), i.kind, i.insight_id))
return tuple(insights), tuple(used), tuple(unavailable)
def load_insights_snapshot(
*,
traffic: Any | None = None,
health: Any | None = None,
provider_snapshot: ProviderSnapshot | None = None,
analytics: Any | None = None,
load_live: bool = True,
) -> InsightsSnapshot:
"""Compose insights from injected or live console evidence sources."""
sources_unavailable: list[dict[str, str]] = []
if load_live and traffic is None and not _offline_test_mode():
try:
from webui.traffic_loader import load_traffic_snapshot
traffic = load_traffic_snapshot()
except Exception as exc: # fail soft
sources_unavailable.append(
{
"source": "traffic",
"reason": str(_redact(f"traffic load failed: {exc}")),
}
)
traffic = None
if load_live and health is None and not _offline_test_mode():
try:
from webui.system_health import load_system_health
health = load_system_health()
except Exception as exc:
sources_unavailable.append(
{
"source": "system_health",
"reason": str(_redact(f"system health load failed: {exc}")),
}
)
health = None
if provider_snapshot is None:
provider_snapshot = load_provider_snapshot()
if load_live and analytics is None and not _offline_test_mode():
try:
from webui.analytics_loader import load_analytics
analytics = load_analytics()
except Exception as exc:
sources_unavailable.append(
{
"source": "analytics",
"reason": str(_redact(f"analytics load failed: {exc}")),
}
)
analytics = None
insights, used, unavailable = generate_insights(
traffic=traffic,
health=health,
provider_snapshot=provider_snapshot,
analytics=analytics,
)
merged_unavailable = tuple(sources_unavailable) + unavailable
# ok when at least one source contributed or we can honestly report absence.
ok = bool(used) or bool(merged_unavailable)
return InsightsSnapshot(
ok=ok,
insights=insights,
sources_used=used,
sources_unavailable=merged_unavailable,
fetch_error=None
if used
else (
"no evidence sources produced a usable snapshot"
if merged_unavailable
else "no insight sources ran"
),
)
def snapshot_providers_to_dict(snapshot: ProviderSnapshot) -> dict[str, Any]:
return snapshot.to_dict()
def snapshot_insights_to_dict(snapshot: InsightsSnapshot) -> dict[str, Any]:
return snapshot.to_dict()
-196
View File
@@ -1,196 +0,0 @@
"""HTML views for AI-provider connections and operational insights (#650)."""
from __future__ import annotations
from html import escape
from webui.insights_loader import (
CONNECTION_DECLARED_AVAILABLE,
CONNECTION_DECLARED_UNAVAILABLE,
InsightsSnapshot,
ProviderSnapshot,
SEVERITY_CRITICAL,
SEVERITY_INFO,
SEVERITY_UNPROVEN,
SEVERITY_WARN,
)
from webui.layout import render_page
_SEVERITY_CSS = {
SEVERITY_CRITICAL: "badge-blocked",
SEVERITY_WARN: "badge-health-degraded",
SEVERITY_INFO: "badge-health-ok",
SEVERITY_UNPROVEN: "badge-health-unproven",
}
_CONN_CSS = {
CONNECTION_DECLARED_AVAILABLE: "badge-health-ok",
CONNECTION_DECLARED_UNAVAILABLE: "badge-health-degraded",
"registry_unavailable": "badge-blocked",
}
def _badge(text: str, css: str) -> str:
return f'<span class="badge {css}">{escape(text)}</span>'
def _limits_card(lines: list[str], *, title: str) -> str:
items = "".join(f"<li>{escape(line)}</li>" for line in lines)
return f"""<div class="prompt-card">
<h3>{escape(title)}</h3>
<ul class="reasons">{items}</ul>
<p class="muted">Advisory surface only — no review, merge, close, or provider
mutation is available here.</p>
</div>"""
def render_providers_page(snapshot: ProviderSnapshot) -> str:
"""Render the AI-provider connections page."""
if not snapshot.ok:
body = f"""<h2>AI provider connections</h2>
<p class="meta">Phase 4 read-only provider status (#650).</p>
<div class="health-card health-stale">
<strong>Provider registry unavailable:</strong>
{escape(snapshot.fetch_error or "registry could not be loaded")}.
No connection table is rendered — an empty table would claim that no
providers are configured.
</div>
{_limits_card([
"connection_status reflects the worker registry declaration only",
"no API keys or credential material are loaded or rendered",
"live executable health is not probed on this surface",
], title="Interpretation limits")}
"""
return render_page(title="Providers", body_html=body)
rows = []
for provider in snapshot.providers:
models = (
", ".join(f"<code>{escape(m)}</code>" for m in provider.models)
if provider.models
else '<span class="muted">none declared</span>'
)
rows.append(
"<tr>"
f"<td><code>{escape(provider.provider_id)}</code></td>"
f"<td>{escape(provider.display_name)}</td>"
f"<td>{escape(provider.vendor)}</td>"
f"<td><code>{escape(provider.executable)}</code></td>"
f"<td>{_badge(provider.connection_status, _CONN_CSS.get(provider.connection_status, 'badge-health-skipped'))}</td>"
f"<td>{provider.worker_count} "
f"({provider.enabled_worker_count} enabled)</td>"
f"<td>{models}</td>"
"</tr>"
)
table = (
"".join(rows)
if rows
else '<tr><td colspan="7" class="muted">No providers declared in the registry.</td></tr>'
)
body = f"""<h2>AI provider connections</h2>
<p class="meta">Phase 4 read-only provider status (#650). Registry revision
<code>{escape(str(snapshot.registry_revision))}</code>.
Declared status only — secrets never load.</p>
<div class="health-card">
<h3>Declared connections</h3>
<table class="registry">
<thead>
<tr>
<th>Provider</th><th>Name</th><th>Vendor</th><th>Executable</th>
<th>Connection</th><th>Workers</th><th>Models (declared)</th>
</tr>
</thead>
<tbody>{table}</tbody>
</table>
<p class="muted">{escape(snapshot.providers[0].probe_limit if snapshot.providers else "")}</p>
</div>
{_limits_card([
"connection_status reflects the worker registry declaration only",
"no API keys or credential material are loaded or rendered",
"live executable health is not probed on this surface (#800 owns that)",
], title="Interpretation limits")}
<p class="muted">Related: <a href="/insights">Operational insights</a> ·
<a href="/analytics">Analytics</a></p>
"""
return render_page(title="Providers", body_html=body)
def _evidence_list(insight) -> str:
items = "".join(
f"<li><code>{escape(ref.kind)}:{escape(ref.ref)}</code> — "
f"{escape(ref.detail)}</li>"
for ref in insight.evidence
)
return f'<ul class="reasons">{items}</ul>'
def render_insights_page(snapshot: InsightsSnapshot) -> str:
"""Render the operational insights page."""
if not snapshot.ok and not snapshot.insights:
body = f"""<h2>Operational insights</h2>
<p class="meta">Phase 4 evidence-backed insights (#650).</p>
<div class="health-card health-stale">
<strong>Insights unavailable:</strong>
{escape(snapshot.fetch_error or "no sources ran")}.
</div>
{_limits_card([
"insights are advisory only and never authorize merge, review, or close",
"an insight without evidence refs is refused rather than emitted",
], title="Interpretation limits")}
"""
return render_page(title="Insights", body_html=body)
source_bits = []
if snapshot.sources_used:
source_bits.append(
"sources used: " + ", ".join(f"<code>{escape(s)}</code>" for s in snapshot.sources_used)
)
if snapshot.sources_unavailable:
missing = "; ".join(
f"{escape(item.get('source', '?'))}: {escape(item.get('reason', ''))}"
for item in snapshot.sources_unavailable
)
source_bits.append(f"sources unavailable: {missing}")
cards = []
for insight in snapshot.insights:
cards.append(
f"""<div class="prompt-card">
<h3>{_badge(insight.severity, _SEVERITY_CSS.get(insight.severity, "badge-health-skipped"))}
{escape(insight.title)}</h3>
<p class="meta"><code>{escape(insight.kind)}</code> · confidence
<code>{escape(insight.confidence)}</code> ·
{_badge("advisory only", "badge-health-skipped")} ·
{_badge("no action claimed", "badge-health-ok")}</p>
<p>{escape(insight.summary)}</p>
<h4>Evidence</h4>
{_evidence_list(insight)}
</div>"""
)
if not cards:
cards.append(
'<div class="prompt-card"><p class="muted">No insights met the '
"evidence threshold in the loaded sources. That is not a claim "
"that the fleet is healthy — only that no qualifying pattern was "
"found.</p></div>"
)
body = f"""<h2>Operational insights</h2>
<p class="meta">Phase 4 evidence-backed insights (#650). Derived only from
durable console evidence; never invents policy or completes workflow actions.</p>
<p class="muted">{" · ".join(source_bits) if source_bits else ""}</p>
{"".join(cards)}
{_limits_card([
"insights are advisory only and never authorize merge, review, or close",
"an insight without evidence refs is refused rather than emitted",
"a missing source is listed as unavailable, not as an empty success",
"insights never claim a workflow action completed",
], title="Interpretation limits")}
<p class="muted">Related: <a href="/providers">Provider connections</a> ·
<a href="/traffic">Traffic</a> · <a href="/system-health">System health</a> ·
<a href="/analytics">Analytics</a></p>
"""
return render_page(title="Insights", body_html=body)
+11 -6
View File
@@ -4,11 +4,12 @@ Single source of truth for the console navigation so ``webui/layout.py`` and
the ``webui/app.py`` route table stay aligned with epic #631. Read-only: every the ``webui/app.py`` route table stay aligned with epic #631. Read-only: every
destination is a GET view or a Phase 1 placeholder. No mutation links. destination is a GET view or a Phase 1 placeholder. No mutation links.
Nav groups follow the #631 information architecture: Health, Traffic, Nav groups follow the #631 Phase 1 information architecture: Health, Traffic,
Runtime/Sessions, Projects, Inventory, Timeline, Policy (placeholder), and Runtime/Sessions, Projects, Inventory, Timeline, Policy (placeholder), and
Gitea linkage (#645) plus Phase 4 Insights/Providers (#650). Later-phase Insights (placeholder), joined by the Phase 3 Gitea linkage group (#645).
surfaces are declared as ``stub`` items and backed by ``STUB_PAGES`` so their Later-phase surfaces are declared as ``stub`` items and
nav links resolve to a graceful placeholder instead of a 404. backed by ``STUB_PAGES`` so their nav links resolve to a graceful placeholder
instead of a 404.
""" """
from __future__ import annotations from __future__ import annotations
@@ -68,8 +69,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
NavItem("/prompts", "Prompts"), NavItem("/prompts", "Prompts"),
)), )),
NavGroup("Insights", ( NavGroup("Insights", (
NavItem("/insights", "Insights"), NavItem("/insights", "Insights", "stub"),
NavItem("/providers", "Providers"),
NavItem("/analytics", "Analytics"), NavItem("/analytics", "Analytics"),
NavItem("/audit", "Audit"), NavItem("/audit", "Audit"),
)), )),
@@ -93,6 +93,11 @@ STUB_PAGES: dict[str, tuple[str, str]] = {
"Policy", "Policy",
"Capability and role policy surface. Placeholder until a later phase.", "Capability and role policy surface. Placeholder until a later phase.",
), ),
"/insights": (
"Insights",
"Aggregate operational insights and trends. Placeholder until a later "
"phase.",
),
} }