Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
af70a27b01 | ||
|
|
983e8ac2c7 |
+24
-98
@@ -23,7 +23,6 @@ import json
|
||||
import os
|
||||
import uuid
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Mapping, Sequence
|
||||
|
||||
from control_plane_db import (
|
||||
@@ -739,46 +738,6 @@ def normalize_exclude_issue_numbers(
|
||||
return sorted(out)
|
||||
|
||||
|
||||
def _claim_expires_at(claim: Any) -> datetime | None:
|
||||
"""Parse a claim's ``expires_at``, or ``None`` when it is absent/malformed."""
|
||||
if not isinstance(claim, Mapping):
|
||||
return None
|
||||
text = str(claim.get("expires_at") or "").strip()
|
||||
if not text:
|
||||
return None
|
||||
if text.endswith("Z"):
|
||||
text = text[:-1] + "+00:00"
|
||||
try:
|
||||
parsed = datetime.fromisoformat(text)
|
||||
except ValueError:
|
||||
return None
|
||||
if parsed.tzinfo is None:
|
||||
parsed = parsed.replace(tzinfo=timezone.utc)
|
||||
return parsed.astimezone(timezone.utc)
|
||||
|
||||
|
||||
def _drop_expired_claims(
|
||||
claims: Mapping[tuple[str, int], dict[str, Any]],
|
||||
*,
|
||||
now: datetime | None = None,
|
||||
) -> dict[tuple[str, int], dict[str, Any]]:
|
||||
"""Claims minus those whose lease has already expired (#643).
|
||||
|
||||
The read-only mirror of ``expire_stale_leases``: the sweep marks such rows
|
||||
``expired`` so they stop being returned as claims, and this reaches the same
|
||||
view without writing. A claim with no parseable ``expires_at`` is **kept** —
|
||||
an unreadable expiry is not evidence that work is free.
|
||||
"""
|
||||
moment = now or datetime.now(timezone.utc)
|
||||
kept: dict[tuple[str, int], dict[str, Any]] = {}
|
||||
for key, claim in (claims or {}).items():
|
||||
expires_at = _claim_expires_at(claim)
|
||||
if expires_at is not None and expires_at <= moment:
|
||||
continue
|
||||
kept[key] = claim
|
||||
return kept
|
||||
|
||||
|
||||
def candidate_set_fingerprint(
|
||||
candidates: Sequence[WorkCandidate],
|
||||
*,
|
||||
@@ -867,22 +826,12 @@ def allocate_next_work(
|
||||
exclude_issue_numbers: Sequence[int] | None = None,
|
||||
expected_candidate_set_fingerprint: str | None = None,
|
||||
allocation_mode: str | None = None,
|
||||
side_effect_free: bool = False,
|
||||
) -> dict[str, Any]:
|
||||
"""Select and optionally reserve the next work unit via control-plane DB.
|
||||
|
||||
*apply=False* (default): dry-run selection only — no lease/assignment.
|
||||
*apply=True*: atomic ``assign_and_lease`` for the selected candidate.
|
||||
|
||||
*side_effect_free* (#643): a dry run that writes **nothing** to the
|
||||
control-plane DB. A plain ``apply=False`` still registered a session row and
|
||||
swept stale leases globally, so a caller advertising a read-only preview was
|
||||
mutating on every call. Under this flag both writes are suppressed and stale
|
||||
leases are instead filtered out of the claim map in memory, which yields the
|
||||
same selection the sweep would have produced without persisting anything.
|
||||
Incompatible with *apply* — the combination fails closed rather than
|
||||
silently reserving.
|
||||
|
||||
*allocation_mode* (#840): ``cross_role`` (default for controller) inspects
|
||||
the complete queue and returns one authoritative selection naming the
|
||||
required downstream role/profile/action. ``role_scoped`` keeps prior
|
||||
@@ -936,57 +885,40 @@ def allocate_next_work(
|
||||
"allocation_mode": (allocation_mode or "").strip() or None,
|
||||
}
|
||||
|
||||
# A side-effect-free run may never reserve: reserving is a write, and the
|
||||
# flag is the caller's assertion that this call writes nothing (#643).
|
||||
if side_effect_free and apply:
|
||||
session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}"
|
||||
try:
|
||||
db.upsert_session(
|
||||
session_id=session_id,
|
||||
role=role_norm,
|
||||
profile=profile_name,
|
||||
pid=os.getpid(),
|
||||
controller_instance_id=controller_instance_id,
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 — surface structured
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"apply": True,
|
||||
"reasons": [
|
||||
"side_effect_free is incompatible with apply=True; an "
|
||||
"assignment is a write (fail closed, #643)"
|
||||
f"failed to register session in control-plane DB: {exc} "
|
||||
"(fail closed, #613)"
|
||||
],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
|
||||
session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}"
|
||||
if not side_effect_free:
|
||||
try:
|
||||
db.upsert_session(
|
||||
session_id=session_id,
|
||||
role=role_norm,
|
||||
profile=profile_name,
|
||||
pid=os.getpid(),
|
||||
controller_instance_id=controller_instance_id,
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 — surface structured
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"reasons": [
|
||||
f"failed to register session in control-plane DB: {exc} "
|
||||
"(fail closed, #613)"
|
||||
],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
|
||||
# Expire stale leases globally before selection.
|
||||
try:
|
||||
db.expire_stale_leases()
|
||||
except Exception as exc: # noqa: BLE001
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"reasons": [f"lease expiry failed: {exc} (fail closed)"],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
# Expire stale leases globally before selection.
|
||||
try:
|
||||
db.expire_stale_leases()
|
||||
except Exception as exc: # noqa: BLE001
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"reasons": [f"lease expiry failed: {exc} (fail closed)"],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
|
||||
terminal = None
|
||||
try:
|
||||
@@ -1021,12 +953,6 @@ def allocate_next_work(
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
if side_effect_free:
|
||||
# ``list_active_claims`` filters on status alone, so without the
|
||||
# global sweep an already-expired lease would still read as a live
|
||||
# claim and the preview would report work as taken that is free.
|
||||
# Drop those in memory: same view the sweep produces, no write.
|
||||
claims = _drop_expired_claims(claims)
|
||||
|
||||
try:
|
||||
exclude_nums = normalize_exclude_issue_numbers(exclude_issue_numbers)
|
||||
|
||||
@@ -1,64 +0,0 @@
|
||||
# Sanctioned Recovery Playbooks & Controls (Phase 2 #644)
|
||||
|
||||
## Overview
|
||||
|
||||
Stale runtimes, worktree binding mismatches, and un-reconciled merged branches previously required expert manual shell recovery. Manual process kills (`pkill -f mcp_server.py`) are strictly forbidden and classified as runtime contamination ([#630](sanctioned-restart-controls.md)).
|
||||
|
||||
Phase 2 introduces **sanctioned recovery playbooks and controls** into the Web Console:
|
||||
- **Diagnose**: Surface stale runtimes, worktree binding errors, contamination markers, and worktree anomalies via health & inventory APIs.
|
||||
- **Preview**: Render mutation ledgers and exact confirmation phrases for recovery playbooks.
|
||||
- **Confirm & Apply**: Execute sanctioned recovery actions through gated, audited paths.
|
||||
- **Verify**: Revalidate control-plane state post-recovery before claiming clean status.
|
||||
|
||||
---
|
||||
|
||||
## Recovery Playbook Taxonomy
|
||||
|
||||
| Playbook ID | Action ID | Minimum Role | Target / Scope | Description |
|
||||
|---|---|---|---|---|
|
||||
| `clear_stale_binding` | `system.clear_stale_binding` | Operator | Active worktree binding | Clear provably missing or superseded `GITEA_ACTIVE_WORKTREE` binding ([#702](../stale_binding_recovery.py)). |
|
||||
| `rebind_session_worktree` | `system.rebind_session_worktree` | Operator | Session worktree | Rebind or synchronize session worktree to verified lease worktree ([#864](../dirty_same_claimant_session_rebind.py)). |
|
||||
| `reconcile_cleanups` | `system.reconcile_cleanups` | Controller | Worktree hygiene | Execute reconciler cleanup preview and apply for merged/superseded PR branches. |
|
||||
| `sanctioned_restart` | `system.restart_namespace` | Admin | MCP Namespace | Restart MCP daemon gracefully via host supervisor ([#642](sanctioned-restart-controls.md)). |
|
||||
|
||||
---
|
||||
|
||||
## Wizard Workflow (Diagnose → Preview → Confirm → Verify)
|
||||
|
||||
### 1. Diagnose (`GET /api/v1/system/recovery/diagnose`)
|
||||
Runs control-plane diagnostics:
|
||||
- **Stale Runtime**: Mismatch between running daemon HEAD, local checkout HEAD, and remote-tracking HEAD.
|
||||
- **Worktree Binding**: Missing path (`provably_stale_missing_path`), unverified inherited binding (`unverified_inherited`), or superseded binding (`superseded_by_session_lease`).
|
||||
- **Contamination**: Checks for live contamination markers from unmanaged process kills.
|
||||
- **Worktree Anomalies**: Scans `branches/` directory for un-reconciled cleanups or missing preserved worktrees.
|
||||
|
||||
Returns `RecoveryDiagnosis` with eligible playbooks.
|
||||
|
||||
### 2. Preview (`POST /api/v1/system/recovery/preview`)
|
||||
Takes `playbook_id` and optional `target`/`params`.
|
||||
Returns:
|
||||
- **Mutation Ledger**: Step-by-step sequence of actions.
|
||||
- **Confirmation Phrase**: Exact phrase required to authorize execution (e.g., `confirm clear_stale_binding`).
|
||||
- **Authorization Decision**: RBAC check against the operator's principal.
|
||||
|
||||
### 3. Apply (`POST /api/v1/system/recovery/apply`)
|
||||
Requires `playbook_id` and matching `confirmation` phrase. Gates run in this order, and each fails closed before anything is mutated:
|
||||
|
||||
1. **RBAC and execution phase** (`console_authz.authorize(..., for_execution=True)`). The phase branch only applies when `for_execution` is set. While `ACTIVE_PHASE` is `1`, every phase-2 recovery action is refused with `phase_not_active`, so no recovery playbook writes yet. Preview reports the same decision under `execution_authorization` / `execution_blocked_reason`.
|
||||
2. **Confirmation phrase** (`confirmation_matches`).
|
||||
3. **Contamination rules** ([#630](sanctioned-restart-controls.md)): the live marker is read from the session inventory and assessed under the gated task key `console_recovery_apply`. A contaminated runtime must be cleared through the reconciler cleanup playbook, which is the one playbook exempted from this gate because it is the designated remedy. The marker is also forwarded to `sanctioned_restart.execute_restart`, so a restart cannot launder a contaminated runtime.
|
||||
|
||||
Apply then executes the sanctioned recovery logic against the **live** process environment — not a copy — and records an audit entry in `console_audit`. A playbook that leaves the binding unchanged reports `performed: false`; `binding_before`, `binding_after`, and `binding_changed` are returned so a no-op cannot read as success.
|
||||
|
||||
Apply does **not** enforce master parity. Parity is reported by Diagnose ([#610](../master_parity_gate.py)) as evidence for the operator; it is not a precondition of this endpoint.
|
||||
|
||||
### 4. Verify (`POST /api/v1/system/recovery/verify`)
|
||||
Re-evaluates control-plane diagnostics post-recovery and **reports** `clean`, `stale_runtime_clean`, `binding_clean`, `binding_classification`, and `contamination_clean`. It reports; it does not assert or block. State is read fresh rather than from the mapping a mutation just wrote. An `unverified_inherited` binding is reported as not clean, because unproven is not clean.
|
||||
|
||||
---
|
||||
|
||||
## Safety & Governance Principles
|
||||
|
||||
1. **No Manual `pkill`**: Direct process killing remains forbidden and is recorded as contamination.
|
||||
2. **Auditability**: Every recovery preview and execution is logged in the console audit trail.
|
||||
3. **Master Parity & Dual Control**: High-privilege recovery actions require controller/admin roles and explicit confirmation phrases.
|
||||
@@ -94,10 +94,6 @@ already define, and a regression test asserts each mapping matches.
|
||||
| `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 |
|
||||
| `system.clear_stale_binding` | operator | gated_write | `gitea.read` | Yes | No | No | 2 |
|
||||
| `system.rebind_session_worktree` | operator | gated_write | `gitea.read` | Yes | No | No | 2 |
|
||||
| `system.reconcile_cleanups` | controller | privileged | `gitea.pr.close` | Yes | No | No | 2 |
|
||||
| `initiate_workflow` | operator | gated_write | `gitea.read` | Yes | No | No | 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
|
||||
@@ -116,12 +112,6 @@ 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).
|
||||
|
||||
`initiate_workflow` (#643) is operator-class because its outcome is a *claim*,
|
||||
not a Gitea verdict. Requesting reviewer or merger work reserves that work
|
||||
through the allocator; it does not grant the right to approve or merge, which
|
||||
stays with the MCP role profile and its own capability gates. See
|
||||
[`webui-requests.md`](webui-requests.md).
|
||||
|
||||
### Authorization decision
|
||||
|
||||
`authorize(action_id, principal, for_execution=False)` returns a decision
|
||||
@@ -136,24 +126,9 @@ record and **denies by default**. The deny reasons are closed and enumerated:
|
||||
| `phase_not_active` | Execution requested for an action whose phase is not open. |
|
||||
| `allowed_preview_only` | Authorized — preview only, execution still disabled. |
|
||||
|
||||
There is no implicit allow branch.
|
||||
|
||||
`execution_enabled` on the decision reports whether the action has a live
|
||||
execution path at all, and is computed by `execution_wired(action)`. There are
|
||||
exactly two ways to be wired:
|
||||
|
||||
1. the action's `phase` is at or below `ACTIVE_PHASE`; or
|
||||
2. the action declares an `execution_env_flag` **and** that variable is set.
|
||||
|
||||
Every action that declares no flag therefore reports `execution_enabled: false`
|
||||
while the console is in Phase 1, so no caller can read an allow as permission
|
||||
to mutate. The per-action flag exists because raising `ACTIVE_PHASE` would
|
||||
enable execution for every action of that phase at once, including ones whose
|
||||
execution path is not implemented. One implemented action goes live on its own
|
||||
flag instead of dragging its unimplemented phase-mates with it.
|
||||
|
||||
`initiate_workflow` is the only action that currently declares a flag
|
||||
(`WEBUI_REQUESTS_EXECUTION`), and it stays denied until an operator sets it.
|
||||
There is no implicit allow branch. Even the allow result reports
|
||||
`execution_enabled: false` while the console is in Phase 1, so no caller can
|
||||
read an allow as permission to mutate.
|
||||
|
||||
## Secret redaction
|
||||
|
||||
@@ -260,22 +235,13 @@ second one. The integration points are already wired and observable:
|
||||
instead of adding a parallel check.
|
||||
- **`GET /api/console/security-model`** publishes the RBAC matrix, redaction
|
||||
policy, and audit policy as JSON for operators and tests.
|
||||
- **`POST /api/v1/requests/preview` and `.../apply`** (#643) are the first
|
||||
actions to use this model for a real execution path. Preview always returns a
|
||||
decision and an audited `previewed` record; apply requires `confirm=true`,
|
||||
emits `succeeded` or `denied`, and reserves work only through the allocator.
|
||||
See [`webui-requests.md`](webui-requests.md).
|
||||
|
||||
A Phase 2 action must: use `execution_wired` rather than a private enable flag,
|
||||
implement the confirmation and dual-control flow the matrix already declares,
|
||||
emit a `succeeded` or `failed` record alongside the `gitea_audit` mutation
|
||||
record, and keep `viewer` unable to reach any of it. Turning on execution
|
||||
without the confirmation flow contradicts a declared requirement and is a
|
||||
review failure, not a shortcut.
|
||||
|
||||
Raising `ACTIVE_PHASE` remains the way to open a whole phase at once, and is
|
||||
deliberately *not* what #643 did: an action-scoped opt-in cannot enable an
|
||||
action whose execution path nobody wrote.
|
||||
To open Phase 2, a child issue must: raise `ACTIVE_PHASE`, implement the
|
||||
confirmation and dual-control flow the matrix already declares, emit a
|
||||
`succeeded` or `failed` record alongside the `gitea_audit` mutation record, and
|
||||
keep `viewer` unable to reach any of it. Turning on execution without the
|
||||
confirmation flow contradicts a declared requirement and is a review failure,
|
||||
not a shortcut.
|
||||
|
||||
## Local-dev mode
|
||||
|
||||
@@ -328,7 +294,6 @@ Until Phase 2 wires it, probe protection rests on network placement alone, as
|
||||
| `WEBUI_ROLE_MAP` | unset | JSON subject → role map |
|
||||
| `WEBUI_REQUIRE_PROBE_AUTH` | unset | Require auth for non-public probes |
|
||||
| `WEBUI_CONSOLE_AUDIT_LOG` | unset | Append-only audit sink path |
|
||||
| `WEBUI_REQUESTS_EXECUTION` | unset | Opt in to `initiate_workflow` execution (#643) |
|
||||
|
||||
All are read server-side only. None is ever rendered into a page or returned by
|
||||
an API.
|
||||
|
||||
+54
-1
@@ -85,7 +85,10 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
|
||||
| `/inventory` | Phase 1 shell stub — unified inventory (backed by #636) |
|
||||
| `/timeline` | Phase 1 shell stub — workflow event timeline |
|
||||
| `/policy` | Phase 1 shell stub — capability/role policy placeholder |
|
||||
| `/insights` | Phase 1 shell stub — operational insights placeholder |
|
||||
| `/providers` | AI-provider connection status (#650) — declared registry only, no secrets |
|
||||
| `/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
|
||||
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
|
||||
@@ -329,6 +332,56 @@ Honesty rules specific to this view:
|
||||
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.
|
||||
|
||||
## 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` is the Phase 3 read-only linkage console: which PR carries which issue,
|
||||
|
||||
@@ -1,160 +0,0 @@
|
||||
# Web console requests: intent preview and workflow initiation (#643)
|
||||
|
||||
**Phase 2. Preview is always live and always read-only. Initiation is wired but
|
||||
denied until an operator opts in.**
|
||||
|
||||
Before this surface, starting role work meant pasting a prompt into a terminal
|
||||
and trusting the operator to have checked the allocator first. Nothing enforced
|
||||
that check, so two sessions could reach for the same issue and each believe it
|
||||
was theirs. This page replaces the paste with a *request*: a desired role, an
|
||||
issue or PR, and a stated intent, answered by an authorization decision and —
|
||||
on confirmation — an exclusive assignment from the allocator.
|
||||
|
||||
| Concern | Module |
|
||||
|---------|--------|
|
||||
| Request model, preview, initiation | `webui/request_service.py` |
|
||||
| Form and preview rendering | `webui/request_views.py` |
|
||||
| Authorization | `webui/console_authz.py` (`initiate_workflow`) |
|
||||
| Audit | `webui/console_audit.py` |
|
||||
| Ownership substrate | `allocator_service.py` + `control_plane_db.py` |
|
||||
|
||||
## Surfaces
|
||||
|
||||
| Path | Method | Purpose |
|
||||
|------|--------|---------|
|
||||
| `/requests` | GET | Request form |
|
||||
| `/requests` | POST | Render an intent preview. **Never assigns.** |
|
||||
| `/api/v1/requests/preview` | POST | Intent preview as JSON |
|
||||
| `/api/v1/requests/apply` | POST | Initiate — confirmed, audited, allocator-owned |
|
||||
|
||||
The HTML form has no initiate button on purpose. Initiating requires a
|
||||
confirmed POST to `/api/v1/requests/apply`, so a stray form submission cannot
|
||||
reserve work as a side effect.
|
||||
|
||||
## The request
|
||||
|
||||
```json
|
||||
{
|
||||
"desired_role": "author",
|
||||
"work_kind": "issue",
|
||||
"work_number": 643,
|
||||
"intent_summary": "implement request preview and initiation",
|
||||
"remote": "prgs",
|
||||
"org": "Scaled-Tech-Consulting",
|
||||
"repo": "Gitea-Tools",
|
||||
"expected_head_sha": null
|
||||
}
|
||||
```
|
||||
|
||||
`desired_role` is one of `author`, `reviewer`, `merger`, `reconciler`,
|
||||
`controller`. `work_kind` is `issue` or `pr`. `remote`/`org`/`repo` default to
|
||||
the first project in the registry when omitted; when neither the request nor
|
||||
the registry resolves them, the request is rejected rather than pointed at some
|
||||
other repository. `intent_summary` is required — it is what the audit record
|
||||
states as the reason — and is truncated to 500 characters.
|
||||
|
||||
Parsing rejects rather than corrects. An unknown role, an unknown work kind, a
|
||||
non-positive number, or a missing intent each return `400` with a `reason_code`
|
||||
and the offending `field`.
|
||||
|
||||
## Preview
|
||||
|
||||
Five checks, each with its own verdict, reason code, and detail:
|
||||
|
||||
| Check | Passes when |
|
||||
|-------|-------------|
|
||||
| `authorization` | The console principal holds `operator` or above |
|
||||
| `capability` | The desired role maps to a declared profile and MCP namespace |
|
||||
| `lease_availability` | No active claim holds the work unit |
|
||||
| `next_safe_action` | The allocator would independently select this exact work unit |
|
||||
| `head_pin` | PR work resolves to a head SHA, and a supplied SHA still matches |
|
||||
|
||||
A preview also returns the role's `allowed_actions` and `prohibited_actions`
|
||||
(from `allocator_service.ROLE_ACTIONS`), the `required_profile` and
|
||||
`required_namespace` the work must run under, and a `correlation_id` that ties
|
||||
the preview to its audit record and to any assignment that follows.
|
||||
|
||||
Preview is read-only in the strict sense: it calls the allocator with
|
||||
`apply=false` and writes nothing but an audit line. An unauthorized principal
|
||||
never reaches the allocator or the control-plane DB at all, so a denial cannot
|
||||
be used to enumerate the queue.
|
||||
|
||||
## Initiation
|
||||
|
||||
`POST /api/v1/requests/apply` refuses in this order, and every refusal returns
|
||||
before any assignment is attempted:
|
||||
|
||||
| Condition | Outcome | Status |
|
||||
|-----------|---------|--------|
|
||||
| Unparseable request | `invalid_request` | 400 |
|
||||
| Not authorized, or execution not wired | `denied` | 403 |
|
||||
| `confirm` not set | `denied` / `confirmation_required` | 409 |
|
||||
| Work unit already claimed | `blocked` / `duplicate_assignment` | 409 |
|
||||
| Allocator would select other work | `wait` / `not_next_safe_work` | 409 |
|
||||
| Allocator declines on apply | `blocked` or `wait` | 409 |
|
||||
| Evidence unavailable | `wait` / `evidence_unavailable` | 503 |
|
||||
| Assigned | `assigned_work` | 201 |
|
||||
|
||||
A success returns the assignment plus a `handoff` block naming the profile, the
|
||||
namespace, and the actions that stay forbidden — enough for the operator to
|
||||
continue in the right MCP namespace without guessing.
|
||||
|
||||
### Why apply runs the allocator twice
|
||||
|
||||
The allocator is the only source of exclusive ownership (#600 / #613), and it
|
||||
selects work; it does not take orders. So `apply` runs a dry-run first and
|
||||
proceeds only when the allocator would independently pick the requested work
|
||||
unit. If it would not, the request reports `wait` and mutates nothing.
|
||||
|
||||
A request is therefore a *confirmation* of the allocator's decision, never an
|
||||
override of it. The apply call carries the dry-run's
|
||||
`candidate_set_fingerprint` as a CAS pin (#776), so a queue that changed
|
||||
between the two calls fails closed rather than assigning against a stale view.
|
||||
The result is checked again on the way out: an assignment naming a different
|
||||
work unit is not read as success.
|
||||
|
||||
### Fail-closed defaults
|
||||
|
||||
- An unreadable control-plane DB denies. It is never treated as "nothing holds
|
||||
this work unit".
|
||||
- An incomplete queue inventory denies (#758). Ranking a partial candidate set
|
||||
can select the wrong work.
|
||||
- An allocator that raises denies.
|
||||
- PR work with no resolvable head SHA denies; a supplied SHA that no longer
|
||||
matches denies with `head_moved`.
|
||||
|
||||
## Enabling initiation
|
||||
|
||||
Execution is wired off. Set `WEBUI_REQUESTS_EXECUTION=1` to enable it for the
|
||||
`initiate_workflow` action only — see
|
||||
[`webui-authz-audit.md`](webui-authz-audit.md) for why this is an
|
||||
action-scoped flag rather than a phase bump. With the variable unset, `apply`
|
||||
returns `403` with `reason_code: unauthorized` no matter who asks.
|
||||
|
||||
Enabling execution does **not** enable approvals or merges. Those are phase 3
|
||||
console actions and remain forbidden in every path here; the console reserves
|
||||
work and hands off, and the MCP role profile enforces what that role may then
|
||||
do.
|
||||
|
||||
## Audit
|
||||
|
||||
Every preview and every apply emits a console audit record (schema in
|
||||
[`webui-authz-audit.md`](webui-authz-audit.md)):
|
||||
|
||||
| Event | `result` |
|
||||
|-------|----------|
|
||||
| Preview | `previewed` |
|
||||
| Refusal at any stage | `denied` |
|
||||
| Assignment created | `succeeded` |
|
||||
|
||||
`correlation.request_id` carries the request's `correlation_id`, and a
|
||||
successful record's `metadata` carries `assignment_id` and `lease_id`, so an
|
||||
assignment can be traced back to the intent that produced it. The operator's
|
||||
`intent_summary` travels in `metadata` and passes through the standard
|
||||
redaction pass before persistence like every other field.
|
||||
|
||||
## Non-goals
|
||||
|
||||
- No browser-initiated approve or merge, in this phase or any other.
|
||||
- No bypass of allocator exclusive ownership; no self-selection of work.
|
||||
- No auto-start from raw monitoring incidents (#612 stays downstream).
|
||||
@@ -63,11 +63,6 @@ CONTAMINATION_GATED_TASKS = frozenset({
|
||||
"merge_pr",
|
||||
"delete_branch",
|
||||
"complete_issue",
|
||||
# Web console recovery playbooks that write (#644). These mutate runtime
|
||||
# binding and process state, so a live contamination marker must block them
|
||||
# exactly as it blocks the Gitea-side mutations above. The reconciler
|
||||
# cleanup playbook is the designated remedy and is exempted by its caller.
|
||||
"console_recovery_apply",
|
||||
})
|
||||
|
||||
CONTAMINATION_KIND = "stable_branch_push"
|
||||
|
||||
@@ -142,23 +142,6 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
|
||||
"permission": "gitea.read",
|
||||
"role": "author",
|
||||
},
|
||||
# #644: Phase 2 Web Console recovery tasks.
|
||||
"clear_stale_binding": {
|
||||
"permission": "gitea.read",
|
||||
"role": "author",
|
||||
},
|
||||
"rebind_session_worktree": {
|
||||
"permission": "gitea.read",
|
||||
"role": "author",
|
||||
},
|
||||
# The console playbook orchestrates gitea_reconcile_merged_cleanups, whose
|
||||
# own gate is gitea.read (matching the existing reconcile_merged_cleanups
|
||||
# entry). Declaring a stricter permission here stated a second, conflicting
|
||||
# authority for one operation.
|
||||
"reconcile_cleanups": {
|
||||
"permission": "gitea.read",
|
||||
"role": "reconciler",
|
||||
},
|
||||
# PR synchronization lifecycle: assess is read-only (any role with gitea.read);
|
||||
# update-by-merge is author-only and mutates the PR head via Gitea API.
|
||||
"assess_pr_sync_status": {
|
||||
|
||||
@@ -7,7 +7,6 @@ import tempfile
|
||||
import threading
|
||||
import unittest
|
||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from allocator_service import (
|
||||
OUTCOME_ASSIGNED,
|
||||
@@ -16,7 +15,6 @@ from allocator_service import (
|
||||
OUTCOME_PREVIEW,
|
||||
OUTCOME_WAIT,
|
||||
WorkCandidate,
|
||||
_drop_expired_claims,
|
||||
allocate_next_work,
|
||||
candidate_from_dict,
|
||||
classify_skip,
|
||||
@@ -364,161 +362,5 @@ class AllocatorServiceTest(unittest.TestCase):
|
||||
self.assertIn("unavailable", res["reasons"][0].lower())
|
||||
|
||||
|
||||
class SideEffectFreeAllocationTest(unittest.TestCase):
|
||||
"""``side_effect_free`` dry runs write nothing to the control plane (#643).
|
||||
|
||||
A plain ``apply=False`` still called ``upsert_session`` and
|
||||
``expire_stale_leases`` before the apply branch was consulted, so a caller
|
||||
advertising a read-only preview mutated on every call — one unreferenced
|
||||
session row per preview, plus a global lease sweep.
|
||||
"""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self._tmp.cleanup()
|
||||
|
||||
def _alloc(self, **kwargs):
|
||||
defaults = dict(
|
||||
db=self.db,
|
||||
session_id="s-preview",
|
||||
role="author",
|
||||
remote="prgs",
|
||||
org="org",
|
||||
repo="repo",
|
||||
candidates=[
|
||||
WorkCandidate(kind="issue", number=643, labels=("status:ready",))
|
||||
],
|
||||
apply=False,
|
||||
profile_name="prgs-author",
|
||||
username="jcwalker3",
|
||||
)
|
||||
defaults.update(kwargs)
|
||||
return allocate_next_work(**defaults)
|
||||
|
||||
def _session_ids(self) -> set[str]:
|
||||
return {str(r.get("session_id")) for r in self.db.list_sessions()}
|
||||
|
||||
def test_side_effect_free_preview_writes_no_session_row(self):
|
||||
before = self._session_ids()
|
||||
result = self._alloc(side_effect_free=True)
|
||||
self.assertEqual(result["outcome"], OUTCOME_PREVIEW)
|
||||
self.assertEqual(self._session_ids(), before)
|
||||
self.assertNotIn("s-preview", self._session_ids())
|
||||
|
||||
def test_plain_dry_run_still_registers_a_session(self):
|
||||
# The default is unchanged for every existing caller.
|
||||
self._alloc()
|
||||
self.assertIn("s-preview", self._session_ids())
|
||||
|
||||
def test_repeated_previews_do_not_accumulate_rows(self):
|
||||
for index in range(5):
|
||||
self._alloc(side_effect_free=True, session_id=f"s-{index}")
|
||||
self.assertEqual(self._session_ids(), set())
|
||||
|
||||
def test_side_effect_free_does_not_sweep_stale_leases(self):
|
||||
self.db.upsert_session(session_id="owner", role="author", pid=1)
|
||||
assigned = self.db.assign_and_lease(
|
||||
session_id="owner",
|
||||
role="author",
|
||||
remote="prgs",
|
||||
org="org",
|
||||
repo="repo",
|
||||
kind="issue",
|
||||
number=999,
|
||||
lease_ttl_seconds=-60, # already expired
|
||||
)
|
||||
self.assertEqual(assigned.outcome, "assigned")
|
||||
|
||||
self._alloc(side_effect_free=True)
|
||||
|
||||
# The expired row is still 'active' in the DB: nothing swept it.
|
||||
statuses = {
|
||||
r["lease_id"]: r["status"]
|
||||
for r in self.db.list_leases(
|
||||
remote="prgs", org="org", repo="repo",
|
||||
statuses=("active", "expired"),
|
||||
)
|
||||
}
|
||||
self.assertEqual(statuses.get(assigned.lease_id), "active")
|
||||
|
||||
def test_expired_claims_are_filtered_in_memory_so_work_stays_selectable(self):
|
||||
"""The read-only mirror of the sweep: expired claims must not block."""
|
||||
self.db.upsert_session(session_id="owner", role="author", pid=1)
|
||||
self.db.assign_and_lease(
|
||||
session_id="owner",
|
||||
role="author",
|
||||
remote="prgs",
|
||||
org="org",
|
||||
repo="repo",
|
||||
kind="issue",
|
||||
number=643,
|
||||
lease_ttl_seconds=-60, # expired: must not withhold #643
|
||||
)
|
||||
result = self._alloc(side_effect_free=True)
|
||||
self.assertEqual(result["outcome"], OUTCOME_PREVIEW)
|
||||
self.assertEqual(result["selected"]["number"], 643)
|
||||
|
||||
def test_a_live_claim_still_withholds_the_work(self):
|
||||
self.db.upsert_session(session_id="owner", role="author", pid=1)
|
||||
self.db.assign_and_lease(
|
||||
session_id="owner",
|
||||
role="author",
|
||||
remote="prgs",
|
||||
org="org",
|
||||
repo="repo",
|
||||
kind="issue",
|
||||
number=643,
|
||||
lease_ttl_seconds=3600,
|
||||
)
|
||||
result = self._alloc(side_effect_free=True)
|
||||
self.assertNotEqual(result["outcome"], OUTCOME_ASSIGNED)
|
||||
self.assertNotEqual((result.get("selected") or {}).get("number"), 643)
|
||||
|
||||
def test_side_effect_free_with_apply_fails_closed(self):
|
||||
result = self._alloc(side_effect_free=True, apply=True)
|
||||
self.assertFalse(result["success"])
|
||||
self.assertEqual(result["outcome"], OUTCOME_NO_SAFE)
|
||||
self.assertIsNone(result["assignment"])
|
||||
self.assertIn("incompatible with apply", result["reasons"][0])
|
||||
# And it reserved nothing.
|
||||
self.assertEqual(
|
||||
self.db.list_leases(remote="prgs", org="org", repo="repo"), []
|
||||
)
|
||||
|
||||
|
||||
class DropExpiredClaimsTest(unittest.TestCase):
|
||||
"""The in-memory expiry filter behind side-effect-free previews (#643)."""
|
||||
|
||||
def test_unparseable_expiry_is_kept_rather_than_assumed_free(self):
|
||||
claims = {
|
||||
("issue", 1): {"lease_id": "l1", "expires_at": "not-a-date"},
|
||||
("issue", 2): {"lease_id": "l2"},
|
||||
("issue", 3): {"lease_id": "l3", "expires_at": None},
|
||||
}
|
||||
self.assertEqual(_drop_expired_claims(claims), claims)
|
||||
|
||||
def test_expired_dropped_and_future_kept(self):
|
||||
now = datetime(2026, 7, 25, 12, 0, tzinfo=timezone.utc)
|
||||
claims = {
|
||||
("issue", 1): {"expires_at": "2026-07-25T11:59:59+00:00"},
|
||||
("issue", 2): {"expires_at": "2026-07-25T12:00:01+00:00"},
|
||||
("issue", 3): {"expires_at": "2026-07-25T12:00:00+00:00"}, # boundary
|
||||
}
|
||||
kept = _drop_expired_claims(claims, now=now)
|
||||
self.assertEqual(set(kept), {("issue", 2)})
|
||||
|
||||
def test_naive_and_zulu_timestamps_are_treated_as_utc(self):
|
||||
now = datetime(2026, 7, 25, 12, 0, tzinfo=timezone.utc)
|
||||
claims = {
|
||||
("issue", 1): {"expires_at": "2026-07-25T11:00:00"}, # naive, past
|
||||
("issue", 2): {"expires_at": "2026-07-25T13:00:00Z"}, # zulu, future
|
||||
}
|
||||
kept = _drop_expired_claims(claims, now=now)
|
||||
self.assertEqual(set(kept), {("issue", 2)})
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -1,506 +0,0 @@
|
||||
"""Unit and integration tests for Phase 2 Web Console recovery controls (#644)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import sys
|
||||
import types
|
||||
import unittest
|
||||
from unittest.mock import patch
|
||||
|
||||
from starlette.testclient import TestClient
|
||||
|
||||
import merged_cleanup_reconcile
|
||||
import runtime_recovery_guard
|
||||
import stable_branch_push_guard
|
||||
import stale_binding_recovery
|
||||
from webui import console_authz, console_recovery, system_health
|
||||
from webui.app import create_app
|
||||
|
||||
|
||||
class TestConsoleRecovery(unittest.TestCase):
|
||||
|
||||
def test_diagnose_recovery_healthy(self) -> None:
|
||||
diag = console_recovery.diagnose_recovery()
|
||||
self.assertIn(diag.status, {console_recovery.STATUS_HEALTHY, console_recovery.STATUS_ACTION_REQUIRED})
|
||||
self.assertIsInstance(diag.playbooks, tuple)
|
||||
self.assertGreaterEqual(len(diag.playbooks), 4)
|
||||
|
||||
playbook_ids = {pb.playbook_id for pb in diag.playbooks}
|
||||
self.assertIn(console_recovery.PLAYBOOK_CLEAR_STALE_BINDING, playbook_ids)
|
||||
self.assertIn(console_recovery.PLAYBOOK_REBIND_SESSION, playbook_ids)
|
||||
self.assertIn(console_recovery.PLAYBOOK_RECONCILE_CLEANUPS, playbook_ids)
|
||||
self.assertIn(console_recovery.PLAYBOOK_SANCTIONED_RESTART, playbook_ids)
|
||||
|
||||
def test_confirmation_phrase_generation_and_matching(self) -> None:
|
||||
phrase = console_recovery.confirmation_phrase("clear_stale_binding")
|
||||
self.assertEqual(phrase, "confirm clear_stale_binding")
|
||||
self.assertTrue(console_recovery.confirmation_matches("clear_stale_binding", "confirm clear_stale_binding"))
|
||||
self.assertFalse(console_recovery.confirmation_matches("clear_stale_binding", "wrong phrase"))
|
||||
|
||||
phrase_target = console_recovery.confirmation_phrase("sanctioned_restart", "gitea-author")
|
||||
self.assertEqual(phrase_target, "confirm sanctioned_restart gitea-author")
|
||||
self.assertTrue(console_recovery.confirmation_matches("sanctioned_restart", "confirm sanctioned_restart gitea-author", "gitea-author"))
|
||||
|
||||
def test_build_recovery_preview(self) -> None:
|
||||
principal = console_authz.Principal("[email protected]", console_authz.OPERATOR, console_authz.IDENTITY_LOCAL_DEV, True)
|
||||
preview = console_recovery.build_recovery_preview(
|
||||
playbook_id=console_recovery.PLAYBOOK_CLEAR_STALE_BINDING,
|
||||
target="test-worktree",
|
||||
principal=principal,
|
||||
)
|
||||
self.assertEqual(preview["playbook_id"], console_recovery.PLAYBOOK_CLEAR_STALE_BINDING)
|
||||
self.assertEqual(preview["action_id"], console_recovery.ACTION_CLEAR_STALE_BINDING)
|
||||
self.assertEqual(preview["confirmation_phrase"], "confirm clear_stale_binding test-worktree")
|
||||
self.assertTrue(len(preview["mutation_ledger"]) >= 3)
|
||||
self.assertTrue(preview["authorization"]["allowed"])
|
||||
|
||||
def test_build_recovery_preview_unknown_playbook(self) -> None:
|
||||
preview = console_recovery.build_recovery_preview("unknown_playbook")
|
||||
self.assertFalse(preview.get("allowed"))
|
||||
self.assertEqual(preview.get("error"), "unknown_playbook")
|
||||
|
||||
def test_execute_recovery_playbook_confirmation_mismatch(self) -> None:
|
||||
# Authorization is checked before confirmation, so the phase gate has to
|
||||
# pass for this test to reach the branch it is about.
|
||||
principal = console_authz.Principal("[email protected]", console_authz.OPERATOR, console_authz.IDENTITY_LOCAL_DEV, True)
|
||||
with self._phase_two_enabled():
|
||||
result = console_recovery.execute_recovery_playbook(
|
||||
playbook_id=console_recovery.PLAYBOOK_CLEAR_STALE_BINDING,
|
||||
confirmation="invalid confirmation",
|
||||
principal=principal,
|
||||
)
|
||||
self.assertFalse(result["success"])
|
||||
self.assertFalse(result["allowed"])
|
||||
self.assertEqual(result["error"], "confirmation_mismatch")
|
||||
|
||||
def test_execute_recovery_playbook_unauthorized(self) -> None:
|
||||
# Anonymous principal has viewer role -> should be denied
|
||||
result = console_recovery.execute_recovery_playbook(
|
||||
playbook_id=console_recovery.PLAYBOOK_CLEAR_STALE_BINDING,
|
||||
confirmation="confirm clear_stale_binding",
|
||||
principal=console_authz.ANONYMOUS,
|
||||
)
|
||||
self.assertFalse(result["success"])
|
||||
self.assertFalse(result["allowed"])
|
||||
self.assertEqual(result["error"], console_authz.DENY_UNAUTHENTICATED)
|
||||
|
||||
def test_execute_refuses_phase_two_write_while_console_is_phase_one(self) -> None:
|
||||
"""B1: the apply path must arm the phase gate, not skip it.
|
||||
|
||||
``authorize`` only applies the phase branch when ``for_execution=True``.
|
||||
The apply path used the default, so an operator executed a phase-2 write
|
||||
while ``ACTIVE_PHASE`` was 1.
|
||||
"""
|
||||
self.assertGreater(
|
||||
console_authz.get_action(console_recovery.ACTION_CLEAR_STALE_BINDING).phase,
|
||||
console_authz.ACTIVE_PHASE,
|
||||
"fixture assumes the recovery actions are ahead of the active phase",
|
||||
)
|
||||
principal = console_authz.Principal(
|
||||
"[email protected]", console_authz.OPERATOR, console_authz.IDENTITY_LOCAL_DEV, True
|
||||
)
|
||||
phrase = console_recovery.confirmation_phrase(
|
||||
console_recovery.PLAYBOOK_CLEAR_STALE_BINDING
|
||||
)
|
||||
result = console_recovery.execute_recovery_playbook(
|
||||
playbook_id=console_recovery.PLAYBOOK_CLEAR_STALE_BINDING,
|
||||
confirmation=phrase,
|
||||
principal=principal,
|
||||
)
|
||||
self.assertFalse(result["success"])
|
||||
self.assertFalse(result["allowed"])
|
||||
self.assertEqual(result["error"], console_authz.DENY_PHASE_NOT_ACTIVE)
|
||||
|
||||
def test_preview_execution_enabled_matches_the_execution_decision(self) -> None:
|
||||
"""B1: preview must not report a bare False it cannot explain."""
|
||||
principal = console_authz.Principal(
|
||||
"[email protected]", console_authz.OPERATOR, console_authz.IDENTITY_LOCAL_DEV, True
|
||||
)
|
||||
preview = console_recovery.build_recovery_preview(
|
||||
playbook_id=console_recovery.PLAYBOOK_REBIND_SESSION,
|
||||
target="branches/feat-issue-644",
|
||||
principal=principal,
|
||||
)
|
||||
self.assertFalse(preview["execution_enabled"])
|
||||
self.assertEqual(
|
||||
preview["execution_blocked_reason"], console_authz.DENY_PHASE_NOT_ACTIVE
|
||||
)
|
||||
self.assertFalse(preview["execution_authorization"]["allowed"])
|
||||
# The preview (non-execution) decision still allows, by role.
|
||||
self.assertTrue(preview["authorization"]["allowed"])
|
||||
|
||||
def _phase_two_enabled(self):
|
||||
"""Raise ACTIVE_PHASE so the execution branches are reachable in tests."""
|
||||
return patch.object(console_authz, "ACTIVE_PHASE", 2)
|
||||
|
||||
def _operator(self) -> console_authz.Principal:
|
||||
return console_authz.Principal(
|
||||
"[email protected]", console_authz.OPERATOR, console_authz.IDENTITY_LOCAL_DEV, True
|
||||
)
|
||||
|
||||
def test_rebind_mutates_the_live_environment_not_a_copy(self) -> None:
|
||||
"""B2: the playbook must change the mapping it claims to have changed."""
|
||||
live_env = {stale_binding_recovery.ACTIVE_WORKTREE_ENV: "branches/stale-old"}
|
||||
phrase = console_recovery.confirmation_phrase(
|
||||
console_recovery.PLAYBOOK_REBIND_SESSION, "branches/feat-issue-644"
|
||||
)
|
||||
with self._phase_two_enabled():
|
||||
result = console_recovery.execute_recovery_playbook(
|
||||
playbook_id=console_recovery.PLAYBOOK_REBIND_SESSION,
|
||||
confirmation=phrase,
|
||||
target="branches/feat-issue-644",
|
||||
principal=self._operator(),
|
||||
env=live_env,
|
||||
)
|
||||
self.assertTrue(result["success"])
|
||||
self.assertEqual(
|
||||
live_env[stale_binding_recovery.ACTIVE_WORKTREE_ENV],
|
||||
"branches/feat-issue-644",
|
||||
"rebind reported success without changing the caller's environment",
|
||||
)
|
||||
self.assertTrue(result["applied_result"]["binding_changed"])
|
||||
self.assertEqual(result["applied_result"]["binding_before"], "branches/stale-old")
|
||||
self.assertEqual(
|
||||
result["applied_result"]["binding_after"], "branches/feat-issue-644"
|
||||
)
|
||||
|
||||
def test_clear_stale_binding_reports_failure_when_nothing_changed(self) -> None:
|
||||
"""B2: a no-op recovery must never be reported as success."""
|
||||
live_env: dict[str, str] = {}
|
||||
phrase = console_recovery.confirmation_phrase(
|
||||
console_recovery.PLAYBOOK_CLEAR_STALE_BINDING
|
||||
)
|
||||
with self._phase_two_enabled():
|
||||
result = console_recovery.execute_recovery_playbook(
|
||||
playbook_id=console_recovery.PLAYBOOK_CLEAR_STALE_BINDING,
|
||||
confirmation=phrase,
|
||||
principal=self._operator(),
|
||||
env=live_env,
|
||||
)
|
||||
self.assertFalse(
|
||||
result["success"],
|
||||
"a clear that changed no binding must not report success",
|
||||
)
|
||||
self.assertFalse(result["applied_result"]["binding_changed"])
|
||||
|
||||
def test_clear_stale_binding_clears_the_live_binding(self) -> None:
|
||||
"""B2: the sanctioned clear must reach the caller's environment."""
|
||||
missing = "/nonexistent/branches/deleted-worktree"
|
||||
live_env = {stale_binding_recovery.ACTIVE_WORKTREE_ENV: missing}
|
||||
phrase = console_recovery.confirmation_phrase(
|
||||
console_recovery.PLAYBOOK_CLEAR_STALE_BINDING
|
||||
)
|
||||
with self._phase_two_enabled():
|
||||
result = console_recovery.execute_recovery_playbook(
|
||||
playbook_id=console_recovery.PLAYBOOK_CLEAR_STALE_BINDING,
|
||||
confirmation=phrase,
|
||||
principal=self._operator(),
|
||||
env=live_env,
|
||||
)
|
||||
if result["success"]:
|
||||
self.assertNotIn(stale_binding_recovery.ACTIVE_WORKTREE_ENV, live_env)
|
||||
self.assertEqual(result["applied_result"]["binding_before"], missing)
|
||||
self.assertIsNone(result["applied_result"]["binding_after"])
|
||||
else:
|
||||
# Fail closed is acceptable; reporting a clear that did not happen
|
||||
# is not. This is the invariant the blocker was about.
|
||||
self.assertFalse(result["applied_result"]["binding_changed"])
|
||||
self.assertEqual(
|
||||
live_env.get(stale_binding_recovery.ACTIVE_WORKTREE_ENV), missing
|
||||
)
|
||||
|
||||
def test_reconcile_playbook_calls_an_entry_point_that_exists(self) -> None:
|
||||
"""B3: the previous call named a function absent from the module."""
|
||||
phrase = console_recovery.confirmation_phrase(
|
||||
console_recovery.PLAYBOOK_RECONCILE_CLEANUPS
|
||||
)
|
||||
fake_server = types.SimpleNamespace(
|
||||
gitea_reconcile_merged_cleanups=lambda **kwargs: {
|
||||
"success": True,
|
||||
"entries": [{"issue_number": 100}],
|
||||
}
|
||||
)
|
||||
with self._phase_two_enabled(), patch.dict(
|
||||
sys.modules, {"gitea_mcp_server": fake_server}
|
||||
):
|
||||
result = console_recovery.execute_recovery_playbook(
|
||||
playbook_id=console_recovery.PLAYBOOK_RECONCILE_CLEANUPS,
|
||||
confirmation=phrase,
|
||||
principal=console_authz.Principal(
|
||||
"[email protected]",
|
||||
console_authz.ADMIN,
|
||||
console_authz.IDENTITY_LOCAL_DEV,
|
||||
True,
|
||||
),
|
||||
)
|
||||
self.assertTrue(result["success"], result.get("applied_result"))
|
||||
self.assertNotIn("error_type", result["applied_result"])
|
||||
self.assertEqual(result["applied_result"]["reconciled_count"], 1)
|
||||
|
||||
def test_reconcile_entry_point_exists_on_the_real_module(self) -> None:
|
||||
"""B3 regression: guard the symbol itself, not just the call shape."""
|
||||
import gitea_mcp_server
|
||||
|
||||
self.assertTrue(
|
||||
hasattr(gitea_mcp_server, "gitea_reconcile_merged_cleanups"),
|
||||
"console recovery depends on this reconciler entry point",
|
||||
)
|
||||
self.assertFalse(
|
||||
hasattr(merged_cleanup_reconcile, "reconcile_merged_cleanups"),
|
||||
"if this module grows the orchestrator, point the playbook back at it",
|
||||
)
|
||||
|
||||
def test_contamination_gate_blocks_a_writing_playbook(self) -> None:
|
||||
"""B4: a live marker plus a gated task key must actually block."""
|
||||
marker = {
|
||||
"kind": "manual_daemon_kill",
|
||||
"reason_class": "manual_daemon_kill",
|
||||
"command_summary": "pkill -f gitea_mcp_server",
|
||||
"active": True,
|
||||
}
|
||||
phrase = console_recovery.confirmation_phrase(
|
||||
console_recovery.PLAYBOOK_REBIND_SESSION, "branches/feat-issue-644"
|
||||
)
|
||||
live_env = {stale_binding_recovery.ACTIVE_WORKTREE_ENV: "branches/stale-old"}
|
||||
with self._phase_two_enabled(), patch.object(
|
||||
console_recovery, "load_active_contamination_marker", return_value=marker
|
||||
):
|
||||
result = console_recovery.execute_recovery_playbook(
|
||||
playbook_id=console_recovery.PLAYBOOK_REBIND_SESSION,
|
||||
confirmation=phrase,
|
||||
target="branches/feat-issue-644",
|
||||
principal=self._operator(),
|
||||
env=live_env,
|
||||
)
|
||||
self.assertFalse(result["success"])
|
||||
self.assertEqual(result["error"], "contaminated_runtime")
|
||||
self.assertEqual(
|
||||
live_env[stale_binding_recovery.ACTIVE_WORKTREE_ENV],
|
||||
"branches/stale-old",
|
||||
"a blocked playbook must not have mutated anything",
|
||||
)
|
||||
|
||||
def test_contamination_gate_exempts_the_reconciler_remedy(self) -> None:
|
||||
"""B4: the designated remedy must stay reachable while contaminated."""
|
||||
marker = {
|
||||
"kind": "manual_daemon_kill",
|
||||
"reason_class": "manual_daemon_kill",
|
||||
"command_summary": "pkill -f gitea_mcp_server",
|
||||
"active": True,
|
||||
}
|
||||
phrase = console_recovery.confirmation_phrase(
|
||||
console_recovery.PLAYBOOK_RECONCILE_CLEANUPS
|
||||
)
|
||||
fake_server = types.SimpleNamespace(
|
||||
gitea_reconcile_merged_cleanups=lambda **kwargs: {
|
||||
"success": True,
|
||||
"entries": [],
|
||||
}
|
||||
)
|
||||
with self._phase_two_enabled(), patch.object(
|
||||
console_recovery, "load_active_contamination_marker", return_value=marker
|
||||
), patch.dict(sys.modules, {"gitea_mcp_server": fake_server}):
|
||||
result = console_recovery.execute_recovery_playbook(
|
||||
playbook_id=console_recovery.PLAYBOOK_RECONCILE_CLEANUPS,
|
||||
confirmation=phrase,
|
||||
principal=console_authz.Principal(
|
||||
"[email protected]",
|
||||
console_authz.ADMIN,
|
||||
console_authz.IDENTITY_LOCAL_DEV,
|
||||
True,
|
||||
),
|
||||
)
|
||||
self.assertNotEqual(result.get("error"), "contaminated_runtime")
|
||||
|
||||
def test_gated_task_key_is_actually_gated(self) -> None:
|
||||
"""B4: the console action id was never a member of the gated set."""
|
||||
self.assertIn(
|
||||
console_recovery.CONTAMINATION_GATED_TASK,
|
||||
stable_branch_push_guard.CONTAMINATION_GATED_TASKS,
|
||||
)
|
||||
self.assertNotIn(
|
||||
console_recovery.ACTION_CLEAR_STALE_BINDING,
|
||||
stable_branch_push_guard.CONTAMINATION_GATED_TASKS,
|
||||
)
|
||||
|
||||
def test_diagnosis_reads_the_key_the_gate_returns(self) -> None:
|
||||
"""B4: ``contaminated`` is a key assess_contamination_gate never returns."""
|
||||
gate = runtime_recovery_guard.assess_contamination_gate(
|
||||
None, task=console_recovery.CONTAMINATION_GATED_TASK, actual_role="operator"
|
||||
)
|
||||
self.assertNotIn("contaminated", gate)
|
||||
self.assertIn("block", gate)
|
||||
|
||||
def test_contaminated_runtime_is_reported_unclean(self) -> None:
|
||||
"""B4: verify_post_recovery reported contamination_clean unconditionally."""
|
||||
marker = {
|
||||
"kind": "manual_daemon_kill",
|
||||
"reason_class": "manual_daemon_kill",
|
||||
"command_summary": "pkill -f gitea_mcp_server",
|
||||
"active": True,
|
||||
}
|
||||
with patch.object(
|
||||
console_recovery, "load_active_contamination_marker", return_value=marker
|
||||
):
|
||||
verification = console_recovery.verify_post_recovery()
|
||||
diag = console_recovery.diagnose_recovery()
|
||||
self.assertFalse(verification["contamination_clean"])
|
||||
self.assertFalse(verification["clean"])
|
||||
self.assertEqual(diag.status, console_recovery.STATUS_BLOCKED_CONTAMINATION)
|
||||
|
||||
def test_master_parity_baseline_is_not_the_head_it_is_compared_against(self) -> None:
|
||||
"""B5: capture_startup_parity was fed the head it was then compared to."""
|
||||
stale = system_health.StaleRuntime(
|
||||
daemon_head="a" * 40,
|
||||
checkout_head="b" * 40,
|
||||
remote_head="b" * 40,
|
||||
stale=True,
|
||||
determinable=True,
|
||||
mutation_safe=False,
|
||||
reasons=("daemon is behind the checkout",),
|
||||
)
|
||||
with patch.object(system_health, "assess_stale_runtime", return_value=stale):
|
||||
diag = console_recovery.diagnose_recovery()
|
||||
parity = diag.master_parity
|
||||
self.assertEqual(parity["startup_head"], "a" * 40)
|
||||
self.assertEqual(parity["current_head"], "b" * 40)
|
||||
self.assertNotEqual(parity["startup_head"], parity["current_head"])
|
||||
self.assertFalse(parity["in_parity"])
|
||||
|
||||
def test_master_parity_carries_the_live_remote_dimension(self) -> None:
|
||||
"""B5: live_remote_head was never passed, dropping the #610 dimension."""
|
||||
stale = system_health.StaleRuntime(
|
||||
daemon_head="c" * 40,
|
||||
checkout_head="c" * 40,
|
||||
remote_head="d" * 40,
|
||||
stale=False,
|
||||
determinable=True,
|
||||
mutation_safe=False,
|
||||
reasons=(),
|
||||
)
|
||||
with patch.object(system_health, "assess_stale_runtime", return_value=stale):
|
||||
diag = console_recovery.diagnose_recovery()
|
||||
self.assertEqual(diag.master_parity.get("live_remote_head"), "d" * 40)
|
||||
|
||||
def test_verify_post_recovery(self) -> None:
|
||||
verification = console_recovery.verify_post_recovery()
|
||||
self.assertIn("clean", verification)
|
||||
self.assertIn("status", verification)
|
||||
self.assertIn("reasons", verification)
|
||||
|
||||
def test_unverified_inherited_binding_is_not_reported_clean(self) -> None:
|
||||
"""B2: ``not clear_eligible`` also read clean for unproven bindings."""
|
||||
binding = {
|
||||
"classification": stale_binding_recovery.CLASSIFICATION_UNVERIFIED_INHERITED,
|
||||
"clear_eligible": False,
|
||||
}
|
||||
diag = console_recovery.diagnose_recovery()
|
||||
patched = console_recovery.RecoveryDiagnosis(
|
||||
status=diag.status,
|
||||
clean=diag.clean,
|
||||
stale_runtime=diag.stale_runtime,
|
||||
master_parity=diag.master_parity,
|
||||
stale_binding=binding,
|
||||
contamination=diag.contamination,
|
||||
worktree_anomalies=diag.worktree_anomalies,
|
||||
playbooks=diag.playbooks,
|
||||
reasons=diag.reasons,
|
||||
)
|
||||
with patch.object(console_recovery, "diagnose_recovery", return_value=patched):
|
||||
verification = console_recovery.verify_post_recovery()
|
||||
self.assertFalse(verification["binding_clean"])
|
||||
self.assertEqual(
|
||||
verification["binding_classification"],
|
||||
stale_binding_recovery.CLASSIFICATION_UNVERIFIED_INHERITED,
|
||||
)
|
||||
|
||||
|
||||
class TestConsoleRecoveryApi(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.app = create_app()
|
||||
self.client = TestClient(self.app)
|
||||
|
||||
def test_api_recovery_diagnose(self) -> None:
|
||||
res = self.client.get("/api/v1/system/recovery/diagnose")
|
||||
self.assertEqual(res.status_code, 200)
|
||||
data = res.json()
|
||||
self.assertIn("status", data)
|
||||
self.assertIn("clean", data)
|
||||
self.assertIn("playbooks", data)
|
||||
self.assertTrue(len(data["playbooks"]) >= 4)
|
||||
|
||||
def test_api_recovery_preview(self) -> None:
|
||||
res = self.client.post(
|
||||
"/api/v1/system/recovery/preview",
|
||||
json={"playbook_id": "clear_stale_binding", "target": "active"},
|
||||
)
|
||||
self.assertEqual(res.status_code, 200)
|
||||
data = res.json()
|
||||
self.assertEqual(data["playbook_id"], "clear_stale_binding")
|
||||
self.assertEqual(data["confirmation_phrase"], "confirm clear_stale_binding active")
|
||||
self.assertIn("mutation_ledger", data)
|
||||
|
||||
def test_api_recovery_apply_denied_without_auth(self) -> None:
|
||||
res = self.client.post(
|
||||
"/api/v1/system/recovery/apply",
|
||||
json={"playbook_id": "clear_stale_binding", "confirmation": "confirm clear_stale_binding"},
|
||||
)
|
||||
self.assertEqual(res.status_code, 400)
|
||||
data = res.json()
|
||||
self.assertFalse(data["success"])
|
||||
self.assertFalse(data["allowed"])
|
||||
|
||||
def test_api_recovery_apply_refuses_phase_two_write_with_dev_auth(self) -> None:
|
||||
"""B1: this previously asserted the phase-gate bypass as intended.
|
||||
|
||||
An authenticated operator posting a valid confirmation still must not
|
||||
execute a phase-2 write while the console is in phase 1. The refusal is
|
||||
the contract; a 200 here means the gate is not armed.
|
||||
"""
|
||||
env = {
|
||||
"WEBUI_AUTH_MODE": "local_dev",
|
||||
"WEBUI_DEV_SUBJECT": "[email protected]",
|
||||
"WEBUI_DEV_ROLE": "operator",
|
||||
}
|
||||
before = os.environ.get("GITEA_ACTIVE_WORKTREE")
|
||||
with patch.dict(os.environ, env):
|
||||
res = self.client.post(
|
||||
"/api/v1/system/recovery/apply",
|
||||
json={
|
||||
"playbook_id": "rebind_session_worktree",
|
||||
"target": "branches/feat-issue-644",
|
||||
"confirmation": "confirm rebind_session_worktree branches/feat-issue-644",
|
||||
},
|
||||
)
|
||||
self.assertEqual(res.status_code, 400)
|
||||
data = res.json()
|
||||
self.assertFalse(data["success"])
|
||||
self.assertFalse(data["allowed"])
|
||||
self.assertEqual(data["error"], console_authz.DENY_PHASE_NOT_ACTIVE)
|
||||
self.assertEqual(
|
||||
os.environ.get("GITEA_ACTIVE_WORKTREE"),
|
||||
before,
|
||||
"a refused apply must not have rebound the live process environment",
|
||||
)
|
||||
|
||||
def test_api_recovery_preview_reports_why_execution_is_disabled(self) -> None:
|
||||
res = self.client.post(
|
||||
"/api/v1/system/recovery/preview",
|
||||
json={"playbook_id": "rebind_session_worktree", "target": "active"},
|
||||
)
|
||||
self.assertEqual(res.status_code, 200)
|
||||
data = res.json()
|
||||
self.assertFalse(data["execution_enabled"])
|
||||
self.assertIn("execution_authorization", data)
|
||||
|
||||
def test_api_recovery_verify(self) -> None:
|
||||
res = self.client.get("/api/v1/system/recovery/verify")
|
||||
self.assertEqual(res.status_code, 200)
|
||||
data = res.json()
|
||||
self.assertIn("clean", data)
|
||||
self.assertIn("status", data)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,404 @@
|
||||
"""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()
|
||||
File diff suppressed because it is too large
Load Diff
+39
-201
@@ -53,6 +53,13 @@ from webui.session_loader import (
|
||||
snapshot_to_dict as session_view_snapshot_to_dict,
|
||||
)
|
||||
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 (
|
||||
load_linkage_snapshot,
|
||||
snapshot_to_dict as linkage_snapshot_to_dict,
|
||||
@@ -77,8 +84,6 @@ from webui.system_health import (
|
||||
snapshot_to_dict as system_health_to_dict,
|
||||
)
|
||||
from webui.system_health_views import render_system_health_page
|
||||
from webui import request_service
|
||||
from webui.request_views import render_requests_page
|
||||
|
||||
_READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"})
|
||||
_AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"})
|
||||
@@ -207,84 +212,6 @@ async def system_health(request: Request) -> HTMLResponse:
|
||||
)
|
||||
|
||||
|
||||
async def api_recovery_diagnose(_request: Request) -> JSONResponse:
|
||||
from webui.console_recovery import diagnose_recovery
|
||||
diag = diagnose_recovery()
|
||||
return JSONResponse({
|
||||
"status": diag.status,
|
||||
"clean": diag.clean,
|
||||
"stale_runtime": diag.stale_runtime,
|
||||
"master_parity": diag.master_parity,
|
||||
"stale_binding": diag.stale_binding,
|
||||
"contamination": diag.contamination,
|
||||
"worktree_anomalies": list(diag.worktree_anomalies),
|
||||
"reasons": list(diag.reasons),
|
||||
"playbooks": [
|
||||
{
|
||||
"playbook_id": pb.playbook_id,
|
||||
"label": pb.label,
|
||||
"action_id": pb.action_id,
|
||||
"description": pb.description,
|
||||
"eligible": pb.eligible,
|
||||
"requires_confirmation": pb.requires_confirmation,
|
||||
"reason": pb.reason,
|
||||
"params_schema": pb.params_schema,
|
||||
}
|
||||
for pb in diag.playbooks
|
||||
],
|
||||
})
|
||||
|
||||
|
||||
async def api_recovery_preview(request: Request) -> JSONResponse:
|
||||
from webui.console_recovery import build_recovery_preview
|
||||
body = {}
|
||||
try:
|
||||
body = await request.json()
|
||||
except Exception:
|
||||
pass
|
||||
playbook_id = body.get("playbook_id") or request.query_params.get("playbook_id") or ""
|
||||
target = body.get("target") or request.query_params.get("target")
|
||||
principal = resolve_principal(request.headers)
|
||||
preview = build_recovery_preview(
|
||||
playbook_id=playbook_id,
|
||||
target=target,
|
||||
params=body,
|
||||
principal=principal,
|
||||
)
|
||||
status = 200 if preview.get("playbook_id") else 400
|
||||
return JSONResponse(preview, status_code=status)
|
||||
|
||||
|
||||
async def api_recovery_apply(request: Request) -> JSONResponse:
|
||||
from webui.console_recovery import execute_recovery_playbook
|
||||
body = {}
|
||||
try:
|
||||
body = await request.json()
|
||||
except Exception:
|
||||
pass
|
||||
playbook_id = body.get("playbook_id", "")
|
||||
confirmation = body.get("confirmation", "")
|
||||
target = body.get("target")
|
||||
principal = resolve_principal(request.headers)
|
||||
request_id = getattr(request.state, "request_id", None)
|
||||
result = execute_recovery_playbook(
|
||||
playbook_id=playbook_id,
|
||||
confirmation=confirmation,
|
||||
target=target,
|
||||
params=body,
|
||||
principal=principal,
|
||||
request_id=request_id,
|
||||
)
|
||||
status_code = 200 if result.get("success") else 400
|
||||
return JSONResponse(result, status_code=status_code)
|
||||
|
||||
|
||||
async def api_recovery_verify(_request: Request) -> JSONResponse:
|
||||
from webui.console_recovery import verify_post_recovery
|
||||
verification = verify_post_recovery()
|
||||
return JSONResponse(verification, status_code=200)
|
||||
|
||||
|
||||
async def queue(_request: Request) -> HTMLResponse:
|
||||
snapshot = load_queue_snapshot()
|
||||
return HTMLResponse(render_page(title="Queue", body_html=render_queue_page(snapshot)))
|
||||
@@ -427,6 +354,34 @@ async def api_sessions(_request: Request) -> JSONResponse:
|
||||
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):
|
||||
"""Load one linkage snapshot from the request's scope and focus parameters."""
|
||||
return load_linkage_snapshot(
|
||||
@@ -460,8 +415,6 @@ async def api_v1_gitea_linkage(request: Request) -> JSONResponse:
|
||||
linkage_snapshot_to_dict(snapshot),
|
||||
status_code=200 if snapshot.ok else 502,
|
||||
)
|
||||
|
||||
|
||||
async def _parse_audit_form(request: Request) -> tuple[str, str | None]:
|
||||
if request.method == "GET":
|
||||
return "", None
|
||||
@@ -859,109 +812,6 @@ async def api_v1_analytics_ingest(request: Request) -> JSONResponse:
|
||||
)
|
||||
|
||||
|
||||
def _default_request_scope() -> dict[str, str]:
|
||||
"""Resolve remote/org/repo from the project registry for request forms.
|
||||
|
||||
Returns an empty mapping when the registry cannot be read, which makes
|
||||
``parse_request`` reject a request that did not name its own scope rather
|
||||
than letting it default to some other repository.
|
||||
"""
|
||||
from webui.queue_loader import _host_from_url # host normalisation helper
|
||||
|
||||
registry, error = _load_project_registry()
|
||||
if error is not None or not registry.projects:
|
||||
return {}
|
||||
project = registry.projects[0]
|
||||
host = _host_from_url(project.remote_host)
|
||||
return {
|
||||
"remote": _derive_remote(host),
|
||||
"org": project.gitea_owner or "",
|
||||
"repo": project.repo_name or "",
|
||||
}
|
||||
|
||||
|
||||
async def _request_payload(request: Request) -> dict[str, object]:
|
||||
"""Read a request body as JSON or form-encoded. Never raises."""
|
||||
content_type = (request.headers.get("content-type") or "").lower()
|
||||
if "application/json" in content_type:
|
||||
try:
|
||||
body = await request.json()
|
||||
except Exception:
|
||||
return {}
|
||||
return dict(body) if isinstance(body, dict) else {}
|
||||
try:
|
||||
form = await request.form()
|
||||
except Exception:
|
||||
return {}
|
||||
return {key: form[key] for key in form}
|
||||
|
||||
|
||||
async def requests_page(request: Request) -> HTMLResponse:
|
||||
"""Operator request form and intent preview (#643).
|
||||
|
||||
POST here only ever *previews*. Initiation is a separate confirmed call to
|
||||
``/api/v1/requests/apply`` so that submitting this form cannot reserve
|
||||
work as a side effect.
|
||||
"""
|
||||
submitted: dict[str, object] = {}
|
||||
preview = None
|
||||
error = None
|
||||
if request.method == "POST":
|
||||
submitted = await _request_payload(request)
|
||||
work_request, error = request_service.parse_request(
|
||||
submitted, default_scope=_default_request_scope()
|
||||
)
|
||||
if work_request is not None:
|
||||
preview = request_service.preview_request(
|
||||
work_request,
|
||||
principal=resolve_principal(headers=dict(request.headers)),
|
||||
)
|
||||
return HTMLResponse(
|
||||
render_requests_page(
|
||||
preview=preview, error=error, submitted=submitted
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
async def api_v1_request_preview(request: Request) -> JSONResponse:
|
||||
"""Dry-run authorization and intent preview for a work request (#643)."""
|
||||
payload = await _request_payload(request)
|
||||
work_request, error = request_service.parse_request(
|
||||
payload, default_scope=_default_request_scope()
|
||||
)
|
||||
if work_request is None:
|
||||
return JSONResponse(error.to_dict(), status_code=400)
|
||||
preview = request_service.preview_request(
|
||||
work_request,
|
||||
principal=resolve_principal(headers=dict(request.headers)),
|
||||
)
|
||||
return JSONResponse(
|
||||
preview.to_dict(), status_code=200 if preview.authorized else 403
|
||||
)
|
||||
|
||||
|
||||
async def api_v1_request_apply(request: Request) -> JSONResponse:
|
||||
"""Initiate a previewed work request through the allocator (#643).
|
||||
|
||||
Fail-closed at every step: unauthorized, unconfirmed, not-next-safe, and
|
||||
already-claimed all return without attempting an assignment.
|
||||
"""
|
||||
payload = await _request_payload(request)
|
||||
work_request, error = request_service.parse_request(
|
||||
payload, default_scope=_default_request_scope()
|
||||
)
|
||||
if work_request is None:
|
||||
return JSONResponse(error.to_dict(), status_code=400)
|
||||
confirm = _truthy_flag(str(payload.get("confirm") or ""))
|
||||
result = request_service.apply_request(
|
||||
work_request,
|
||||
principal=resolve_principal(headers=dict(request.headers)),
|
||||
confirm=confirm,
|
||||
)
|
||||
status = int(result.pop("status_code", 403))
|
||||
return JSONResponse(result, status_code=status)
|
||||
|
||||
|
||||
async def method_not_allowed(request: Request, _exc: Exception) -> Response:
|
||||
path = request.url.path
|
||||
if path in _AUDIT_MUTATION_PATHS and request.method == "POST":
|
||||
@@ -1008,6 +858,10 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
|
||||
Route("/api/sessions", api_sessions, methods=["GET"]),
|
||||
Route("/api/v1/sessions", api_sessions, 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("/api/v1/gitea/linkage", api_v1_gitea_linkage, methods=["GET"]),
|
||||
Route("/analytics", analytics, methods=["GET"]),
|
||||
@@ -1031,17 +885,6 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
|
||||
api_action_attempt,
|
||||
methods=["POST"],
|
||||
),
|
||||
Route("/requests", requests_page, methods=["GET", "POST"]),
|
||||
Route(
|
||||
"/api/v1/requests/preview",
|
||||
api_v1_request_preview,
|
||||
methods=["POST"],
|
||||
),
|
||||
Route(
|
||||
"/api/v1/requests/apply",
|
||||
api_v1_request_apply,
|
||||
methods=["POST"],
|
||||
),
|
||||
Route("/api/leases", api_leases, methods=["GET"]),
|
||||
Route("/api/v1/inventory", api_inventory, methods=["GET"]),
|
||||
Route(
|
||||
@@ -1054,11 +897,6 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
|
||||
api_console_security_model,
|
||||
methods=["GET"],
|
||||
),
|
||||
# #644 Phase 2 Recovery API routes
|
||||
Route("/api/v1/system/recovery/diagnose", api_recovery_diagnose, methods=["GET"]),
|
||||
Route("/api/v1/system/recovery/preview", api_recovery_preview, methods=["POST", "GET"]),
|
||||
Route("/api/v1/system/recovery/apply", api_recovery_apply, methods=["POST"]),
|
||||
Route("/api/v1/system/recovery/verify", api_recovery_verify, methods=["POST", "GET"]),
|
||||
*[
|
||||
Route(path, phase_stub, methods=["GET"])
|
||||
for path in STUB_PAGES
|
||||
|
||||
+8
-105
@@ -115,12 +115,6 @@ class ConsoleAction:
|
||||
break_glass: bool
|
||||
phase: int
|
||||
summary: str
|
||||
# Opt-in switch for an action whose execution path is genuinely wired
|
||||
# ahead of its phase becoming globally active (#643). Naming a variable
|
||||
# here enables nothing on its own: the variable must also be set in the
|
||||
# environment. An action that leaves this ``None`` can only execute once
|
||||
# ACTIVE_PHASE reaches its phase, exactly as before.
|
||||
execution_env_flag: str | None = None
|
||||
|
||||
@property
|
||||
def mcp_permission(self) -> str:
|
||||
@@ -283,61 +277,6 @@ _ACTION_SPECS: tuple[ConsoleAction, ...] = (
|
||||
phase=2,
|
||||
summary="Restart one MCP namespace via the host supervisor.",
|
||||
),
|
||||
# #644: Phase 2 recovery controls & playbooks.
|
||||
ConsoleAction(
|
||||
action_id="system.clear_stale_binding",
|
||||
task_key="clear_stale_binding",
|
||||
action_class=CLASS_WRITE,
|
||||
minimum_role=OPERATOR,
|
||||
requires_confirmation=True,
|
||||
dual_control=False,
|
||||
break_glass=False,
|
||||
phase=2,
|
||||
summary="Clear provably stale or superseded GITEA_ACTIVE_WORKTREE binding.",
|
||||
),
|
||||
ConsoleAction(
|
||||
action_id="system.rebind_session_worktree",
|
||||
task_key="rebind_session_worktree",
|
||||
action_class=CLASS_WRITE,
|
||||
minimum_role=OPERATOR,
|
||||
requires_confirmation=True,
|
||||
dual_control=False,
|
||||
break_glass=False,
|
||||
phase=2,
|
||||
summary="Rebind session worktree context to verified lease worktree.",
|
||||
),
|
||||
ConsoleAction(
|
||||
action_id="system.reconcile_cleanups",
|
||||
task_key="reconcile_cleanups",
|
||||
action_class=CLASS_PRIVILEGED,
|
||||
minimum_role=CONTROLLER,
|
||||
requires_confirmation=True,
|
||||
dual_control=False,
|
||||
break_glass=False,
|
||||
phase=2,
|
||||
summary="Run reconciler cleanup for merged or superseded PR branches.",
|
||||
),
|
||||
# #643: submit a work request — desired role, issue/PR, intent — and let
|
||||
# the allocator reserve it. This is the one Phase 2 action whose execution
|
||||
# path is actually implemented (``webui.request_service``), so it carries
|
||||
# the opt-in flag; it stays denied until an operator sets that variable.
|
||||
# Authority is operator-class because the outcome is a claim, not a Gitea
|
||||
# verdict: initiating reviewer or merger *work* does not grant the right
|
||||
# to approve or merge, which stays with the MCP role profile.
|
||||
ConsoleAction(
|
||||
action_id="initiate_workflow",
|
||||
task_key="allocate_next_work",
|
||||
action_class=CLASS_WRITE,
|
||||
minimum_role=OPERATOR,
|
||||
requires_confirmation=True,
|
||||
dual_control=False,
|
||||
break_glass=False,
|
||||
phase=2,
|
||||
summary=(
|
||||
"Preview and initiate allocator-owned workflow work for a role."
|
||||
),
|
||||
execution_env_flag="WEBUI_REQUESTS_EXECUTION",
|
||||
),
|
||||
)
|
||||
|
||||
ACTIONS: dict[str, ConsoleAction] = {a.action_id: a for a in _ACTION_SPECS}
|
||||
@@ -491,33 +430,6 @@ ALLOW_PREVIEW = "allowed_preview_only"
|
||||
# gated on this model landing; nothing here enables it.
|
||||
ACTIVE_PHASE = 1
|
||||
|
||||
_TRUTHY = frozenset({"1", "true", "yes", "on"})
|
||||
|
||||
|
||||
def execution_wired(
|
||||
action: ConsoleAction | None, env: dict[str, str] | None = None
|
||||
) -> bool:
|
||||
"""Whether *action* has a live execution path right now.
|
||||
|
||||
Two ways to be wired, and only two. The action's phase is active, or the
|
||||
action declares an opt-in environment variable *and* that variable is set.
|
||||
Everything else — including every action that never declares a flag — is
|
||||
unwired, so the default across the registry stays deny.
|
||||
|
||||
Bumping ``ACTIVE_PHASE`` would enable execution for every action of that
|
||||
phase at once. The per-action flag exists so a single implemented action
|
||||
can go live without dragging its unimplemented phase-mates with it.
|
||||
"""
|
||||
if action is None:
|
||||
return False
|
||||
if action.phase <= ACTIVE_PHASE:
|
||||
return True
|
||||
flag = (action.execution_env_flag or "").strip()
|
||||
if not flag:
|
||||
return False
|
||||
source = env if env is not None else os.environ
|
||||
return (source.get(flag) or "").strip().lower() in _TRUTHY
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class AuthorizationDecision:
|
||||
@@ -557,19 +469,16 @@ def authorize(
|
||||
principal: Principal | None = None,
|
||||
*,
|
||||
for_execution: bool = False,
|
||||
env: dict[str, str] | None = None,
|
||||
) -> AuthorizationDecision:
|
||||
"""Decide whether *principal* may invoke *action_id*. Deny by default.
|
||||
|
||||
``for_execution`` distinguishes a read-only preview from a real invocation.
|
||||
``execution_enabled`` reports whether the action has a live execution path
|
||||
at all (:func:`execution_wired`) — for every action without an explicit
|
||||
opt-in flag that stays ``False`` while the console is in Phase 1, so no
|
||||
caller can read an allow as permission to mutate.
|
||||
Even an allowed decision reports ``execution_enabled=False`` while the
|
||||
console is in Phase 1, so no caller can read an allow as permission to
|
||||
mutate.
|
||||
"""
|
||||
who = principal if principal is not None else ANONYMOUS
|
||||
action = get_action(action_id)
|
||||
wired = execution_wired(action, env)
|
||||
|
||||
if action is None:
|
||||
return AuthorizationDecision(
|
||||
@@ -588,7 +497,7 @@ def authorize(
|
||||
"requires_confirmation": action.requires_confirmation,
|
||||
"dual_control": action.dual_control,
|
||||
"break_glass": action.break_glass,
|
||||
"execution_enabled": wired,
|
||||
"execution_enabled": False,
|
||||
}
|
||||
|
||||
if not who.authenticated:
|
||||
@@ -621,19 +530,13 @@ def authorize(
|
||||
**base,
|
||||
)
|
||||
|
||||
if for_execution and not wired:
|
||||
if for_execution and action.phase > ACTIVE_PHASE:
|
||||
return AuthorizationDecision(
|
||||
allowed=False,
|
||||
reason_code=DENY_PHASE_NOT_ACTIVE,
|
||||
detail=(
|
||||
f"Action {action_id!r} belongs to phase {action.phase}; the "
|
||||
f"console is in phase {ACTIVE_PHASE}"
|
||||
+ (
|
||||
f" and {action.execution_env_flag} is not set"
|
||||
if action.execution_env_flag
|
||||
else ""
|
||||
)
|
||||
+ ". Execution is not wired."
|
||||
f"console is in phase {ACTIVE_PHASE}. Execution is not wired."
|
||||
),
|
||||
**base,
|
||||
)
|
||||
@@ -642,8 +545,8 @@ def authorize(
|
||||
allowed=True,
|
||||
reason_code=ALLOW_PREVIEW,
|
||||
detail=(
|
||||
"Principal holds the required role. Execution proceeds only for an "
|
||||
"action with a wired execution path; everything else is preview."
|
||||
"Principal holds the required role. Preview only — execution "
|
||||
"remains disabled until the Phase 2 action framework ships."
|
||||
),
|
||||
**base,
|
||||
)
|
||||
|
||||
@@ -1,721 +0,0 @@
|
||||
"""Web Console Phase 2 Recovery Controls & Playbooks (#644).
|
||||
|
||||
Provides canonical recovery controls for the web console:
|
||||
1. Diagnosis: Surfaces stale runtimes, worktree binding errors, contamination markers,
|
||||
and un-reconciled cleanups.
|
||||
2. Gated Actions & Playbooks: Guided recovery (rebind session worktree, clear stale
|
||||
binding, trigger reconciler cleanups, sanctioned restart).
|
||||
3. RBAC, Contamination (#630), and Master Parity (#610) integration.
|
||||
4. Audit trail via ``console_audit`` and mandatory post-recovery revalidation.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from dataclasses import asdict, dataclass, field
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
import master_parity_gate
|
||||
import runtime_recovery_guard
|
||||
import stale_binding_recovery
|
||||
from webui import console_audit, console_authz, sanctioned_restart, system_health, worktree_scanner
|
||||
|
||||
# --- Recovery Statuses ------------------------------------------------------
|
||||
STATUS_HEALTHY = "healthy"
|
||||
STATUS_ACTION_REQUIRED = "action_required"
|
||||
STATUS_BLOCKED_CONTAMINATION = "blocked_contamination"
|
||||
STATUS_RECONNECT_REQUIRED = "reconnect_required"
|
||||
|
||||
# --- Playbook Identifiers ---------------------------------------------------
|
||||
PLAYBOOK_CLEAR_STALE_BINDING = "clear_stale_binding"
|
||||
PLAYBOOK_REBIND_SESSION = "rebind_session_worktree"
|
||||
PLAYBOOK_RECONCILE_CLEANUPS = "reconcile_cleanups"
|
||||
PLAYBOOK_SANCTIONED_RESTART = "sanctioned_restart"
|
||||
|
||||
KNOWN_PLAYBOOKS: tuple[str, ...] = (
|
||||
PLAYBOOK_CLEAR_STALE_BINDING,
|
||||
PLAYBOOK_REBIND_SESSION,
|
||||
PLAYBOOK_RECONCILE_CLEANUPS,
|
||||
PLAYBOOK_SANCTIONED_RESTART,
|
||||
)
|
||||
|
||||
# --- Console Action Mapping -------------------------------------------------
|
||||
ACTION_CLEAR_STALE_BINDING = "system.clear_stale_binding"
|
||||
ACTION_REBIND_SESSION = "system.rebind_session_worktree"
|
||||
ACTION_RECONCILE_CLEANUPS = "system.reconcile_cleanups"
|
||||
|
||||
PLAYBOOK_ACTIONS: dict[str, str] = {
|
||||
PLAYBOOK_CLEAR_STALE_BINDING: ACTION_CLEAR_STALE_BINDING,
|
||||
PLAYBOOK_REBIND_SESSION: ACTION_REBIND_SESSION,
|
||||
PLAYBOOK_RECONCILE_CLEANUPS: ACTION_RECONCILE_CLEANUPS,
|
||||
PLAYBOOK_SANCTIONED_RESTART: sanctioned_restart.ACTION_RESTART_NAMESPACE,
|
||||
}
|
||||
|
||||
#: Task key handed to :func:`runtime_recovery_guard.assess_contamination_gate`.
|
||||
#: A console *action id* is not a task name and is not a member of
|
||||
#: ``CONTAMINATION_GATED_TASKS``, so passing one left the #630 gate inert. Every
|
||||
#: writing recovery playbook shares this one gated task key; the reconciler
|
||||
#: cleanup playbook is exempted separately because it is the designated remedy.
|
||||
CONTAMINATION_GATED_TASK = "console_recovery_apply"
|
||||
|
||||
#: Remote whose contamination markers govern this console. Markers are written
|
||||
#: per remote, so reading the wrong one reports a contaminated runtime clean.
|
||||
REMOTE_ENV = "WEBUI_GITEA_REMOTE"
|
||||
DEFAULT_REMOTE = "prgs"
|
||||
|
||||
|
||||
def _console_remote(env: dict[str, str] | None = None) -> str:
|
||||
env_map = env if env is not None else os.environ
|
||||
return (env_map.get(REMOTE_ENV) or "").strip() or DEFAULT_REMOTE
|
||||
|
||||
|
||||
def load_active_contamination_marker(
|
||||
remote: str | None = None, env: dict[str, str] | None = None
|
||||
) -> dict[str, Any] | None:
|
||||
"""Return the live #630 contamination marker payload, or ``None``.
|
||||
|
||||
The gate is only meaningful when it is fed a real marker: with
|
||||
``marker=None`` :func:`assess_contamination_gate` returns ``block: False``
|
||||
on its first statement. The #641 session inventory already reads the durable
|
||||
markers, so reuse that reader rather than adding a second source of truth.
|
||||
Never raises into a diagnosis or execution path.
|
||||
"""
|
||||
try:
|
||||
from webui import session_loader
|
||||
except Exception: # noqa: BLE001 — never break recovery on an import problem
|
||||
return None
|
||||
try:
|
||||
markers = session_loader._load_contamination_markers(
|
||||
remote=remote or _console_remote(env)
|
||||
)
|
||||
except Exception: # noqa: BLE001 — fail soft; the caller degrades to no marker
|
||||
return None
|
||||
for marker in markers:
|
||||
payload = marker.to_dict()
|
||||
if payload.get("active"):
|
||||
return payload
|
||||
return None
|
||||
|
||||
|
||||
def _active_binding(env_map: Any) -> str | None:
|
||||
"""Read the live worktree binding so a no-op recovery cannot report success."""
|
||||
value = env_map.get(stale_binding_recovery.ACTIVE_WORKTREE_ENV)
|
||||
return value if value else None
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class RecoveryLedgerEntry:
|
||||
"""One planned recovery step displayed before execution."""
|
||||
|
||||
sequence: int
|
||||
step: str
|
||||
summary: str
|
||||
executes_process_kill: bool = False
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class PlaybookDescriptor:
|
||||
"""Structured recovery playbook option returned during diagnosis."""
|
||||
|
||||
playbook_id: str
|
||||
label: str
|
||||
action_id: str
|
||||
description: str
|
||||
eligible: bool
|
||||
requires_confirmation: bool
|
||||
reason: str
|
||||
params_schema: dict[str, Any] = field(default_factory=dict)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class RecoveryDiagnosis:
|
||||
"""Complete diagnostic snapshot of control-plane recovery needs."""
|
||||
|
||||
status: str
|
||||
clean: bool
|
||||
stale_runtime: dict[str, Any]
|
||||
master_parity: dict[str, Any]
|
||||
stale_binding: dict[str, Any]
|
||||
contamination: dict[str, Any]
|
||||
worktree_anomalies: tuple[str, ...]
|
||||
playbooks: tuple[PlaybookDescriptor, ...]
|
||||
reasons: tuple[str, ...]
|
||||
|
||||
|
||||
def _repo_root(custom_path: Path | str | None = None) -> Path:
|
||||
if custom_path:
|
||||
return Path(custom_path).resolve()
|
||||
override = (os.environ.get("WEBUI_REPO_ROOT") or "").strip()
|
||||
if override:
|
||||
return Path(override).resolve()
|
||||
return Path(__file__).resolve().parent.parent
|
||||
|
||||
|
||||
def confirmation_phrase(playbook_id: str, target: str | None = None) -> str:
|
||||
"""Construct exact confirmation phrase required for a recovery playbook."""
|
||||
clean_target = (target or "").strip()
|
||||
if clean_target:
|
||||
return f"confirm {playbook_id} {clean_target}"
|
||||
return f"confirm {playbook_id}"
|
||||
|
||||
|
||||
def confirmation_matches(
|
||||
playbook_id: str, confirmation: str | None, target: str | None = None
|
||||
) -> bool:
|
||||
expected = confirmation_phrase(playbook_id, target)
|
||||
return str(confirmation or "").strip() == expected
|
||||
|
||||
|
||||
def _build_ledger(
|
||||
playbook_id: str, target: str | None = None
|
||||
) -> tuple[RecoveryLedgerEntry, ...]:
|
||||
if playbook_id == PLAYBOOK_CLEAR_STALE_BINDING:
|
||||
return (
|
||||
RecoveryLedgerEntry(1, "quiesce", "Stop admitting new gated mutations."),
|
||||
RecoveryLedgerEntry(
|
||||
2,
|
||||
"clear_env",
|
||||
f"Remove stale env binding GITEA_ACTIVE_WORKTREE ({target or 'active'}).",
|
||||
),
|
||||
RecoveryLedgerEntry(
|
||||
3, "audit", "Record clear_stale_binding event in console audit log."
|
||||
),
|
||||
RecoveryLedgerEntry(
|
||||
4, "revalidate", "Re-run diagnosis to verify clean binding state."
|
||||
),
|
||||
)
|
||||
if playbook_id == PLAYBOOK_REBIND_SESSION:
|
||||
return (
|
||||
RecoveryLedgerEntry(1, "quiesce", "Stop admitting new gated mutations."),
|
||||
RecoveryLedgerEntry(
|
||||
2,
|
||||
"rebind_worktree",
|
||||
f"Rebind session worktree context safely to {target or 'target worktree'}.",
|
||||
),
|
||||
RecoveryLedgerEntry(
|
||||
3, "audit", "Record rebind_session_worktree event in console audit log."
|
||||
),
|
||||
RecoveryLedgerEntry(
|
||||
4, "revalidate", "Re-run diagnosis to verify worktree binding state."
|
||||
),
|
||||
)
|
||||
if playbook_id == PLAYBOOK_RECONCILE_CLEANUPS:
|
||||
return (
|
||||
RecoveryLedgerEntry(1, "quiesce", "Stop admitting new gated mutations."),
|
||||
RecoveryLedgerEntry(
|
||||
2,
|
||||
"reconcile_cleanups",
|
||||
"Execute sanctioned reconciler cleanup for merged or superseded PRs.",
|
||||
),
|
||||
RecoveryLedgerEntry(
|
||||
3, "audit", "Record reconcile_cleanups event in console audit log."
|
||||
),
|
||||
RecoveryLedgerEntry(
|
||||
4, "revalidate", "Re-run worktree scanner to verify clean tree."
|
||||
),
|
||||
)
|
||||
if playbook_id == PLAYBOOK_SANCTIONED_RESTART:
|
||||
restart_ledger = sanctioned_restart._mutation_ledger(
|
||||
target or "gitea-author", sanctioned_restart.MODE_RESTART
|
||||
)
|
||||
return tuple(
|
||||
RecoveryLedgerEntry(
|
||||
sequence=e.sequence,
|
||||
step=e.step,
|
||||
summary=e.summary,
|
||||
executes_process_kill=e.executes_process_kill,
|
||||
)
|
||||
for e in restart_ledger
|
||||
)
|
||||
return (
|
||||
RecoveryLedgerEntry(1, "unspecified", f"Execute recovery playbook {playbook_id}."),
|
||||
)
|
||||
|
||||
|
||||
def diagnose_recovery(
|
||||
repo_path: Path | str | None = None,
|
||||
env: dict[str, str] | None = None,
|
||||
active_worktree_val: str | None = None,
|
||||
session_lease_wt: str | None = None,
|
||||
role_kind: str | None = None,
|
||||
) -> RecoveryDiagnosis:
|
||||
"""Run full control-plane diagnostics to determine recovery needs and options."""
|
||||
root = _repo_root(repo_path)
|
||||
source_env = dict(env) if env is not None else dict(os.environ)
|
||||
reasons: list[str] = []
|
||||
|
||||
# 1. Stale runtime assessment
|
||||
stale_runtime_obj = system_health.assess_stale_runtime(root)
|
||||
stale_runtime_dict = {
|
||||
"daemon_head": stale_runtime_obj.daemon_head,
|
||||
"checkout_head": stale_runtime_obj.checkout_head,
|
||||
"remote_head": stale_runtime_obj.remote_head,
|
||||
"stale": stale_runtime_obj.stale,
|
||||
"determinable": stale_runtime_obj.determinable,
|
||||
"mutation_safe": stale_runtime_obj.mutation_safe,
|
||||
"reasons": list(stale_runtime_obj.reasons),
|
||||
}
|
||||
if stale_runtime_obj.stale:
|
||||
reasons.append("Runtime HEAD disagrees with checkout/remote HEAD.")
|
||||
|
||||
# 2. Master parity assessment
|
||||
#
|
||||
# The baseline is the commit the *running process* started at, which is what
|
||||
# the parity gate is about. Capturing it from ``checkout_head`` and then
|
||||
# comparing it against that same value made ``in_parity`` structurally
|
||||
# incapable of being false. ``live_remote_head`` restores the #610
|
||||
# live-remote dimension, which was previously dropped.
|
||||
checkout_head = stale_runtime_obj.checkout_head
|
||||
startup_dict = master_parity_gate.capture_startup_parity(
|
||||
str(root), head=stale_runtime_obj.daemon_head
|
||||
)
|
||||
parity_dict = master_parity_gate.assess_master_parity(
|
||||
startup_dict, checkout_head, stale_runtime_obj.remote_head
|
||||
)
|
||||
if not parity_dict.get("in_parity", True):
|
||||
reasons.append("Repository is not in master parity.")
|
||||
|
||||
# 3. Worktree binding classification
|
||||
boot_bindings = stale_binding_recovery.snapshot_boot_bindings(source_env)
|
||||
active_val = (
|
||||
active_worktree_val
|
||||
if active_worktree_val is not None
|
||||
else source_env.get(stale_binding_recovery.ACTIVE_WORKTREE_ENV)
|
||||
)
|
||||
boot_inherited = bool(boot_bindings.get("active_worktree") and active_val == boot_bindings.get("active_worktree"))
|
||||
|
||||
path_exists = None
|
||||
if active_val:
|
||||
path_exists = os.path.exists(os.path.realpath(active_val))
|
||||
|
||||
binding_class = stale_binding_recovery.classify_active_worktree_binding(
|
||||
active_value=active_val,
|
||||
session_lease_worktree=session_lease_wt,
|
||||
boot_inherited=boot_inherited,
|
||||
path_exists=path_exists,
|
||||
role_kind=role_kind,
|
||||
)
|
||||
|
||||
if binding_class.get("clear_eligible"):
|
||||
reasons.append(
|
||||
f"Active worktree binding is stale ({binding_class.get('classification')})."
|
||||
)
|
||||
elif binding_class.get("classification") == stale_binding_recovery.CLASSIFICATION_UNVERIFIED_INHERITED:
|
||||
reasons.append("Inherited worktree binding is unverified.")
|
||||
|
||||
# 4. Contamination assessment (#630)
|
||||
#
|
||||
# A real marker and a task key the gate actually gates on: with marker=None
|
||||
# the gate short-circuits to ``block: False``, and with a console action id
|
||||
# the task is outside CONTAMINATION_GATED_TASKS, so it could never block.
|
||||
contamination_marker = load_active_contamination_marker(env=source_env)
|
||||
contamination_dict = runtime_recovery_guard.assess_contamination_gate(
|
||||
contamination_marker,
|
||||
task=CONTAMINATION_GATED_TASK,
|
||||
actual_role=role_kind,
|
||||
)
|
||||
contaminated = bool(contamination_dict.get("block"))
|
||||
if contaminated:
|
||||
reasons.append("Runtime is contaminated by manual process kill (#630).")
|
||||
|
||||
# 5. Worktree scanner hygiene & anomalies
|
||||
hygiene = worktree_scanner.load_hygiene_snapshot(project_root=str(root))
|
||||
worktree_anomalies = hygiene.anomalies
|
||||
|
||||
# Determine status & eligible playbooks
|
||||
playbooks: list[PlaybookDescriptor] = []
|
||||
|
||||
# Playbook 1: Clear Stale Binding
|
||||
clear_eligible = bool(binding_class.get("clear_eligible"))
|
||||
playbooks.append(
|
||||
PlaybookDescriptor(
|
||||
playbook_id=PLAYBOOK_CLEAR_STALE_BINDING,
|
||||
label="Clear Stale Worktree Binding",
|
||||
action_id=ACTION_CLEAR_STALE_BINDING,
|
||||
description="Clear provably stale or superseded GITEA_ACTIVE_WORKTREE environment binding.",
|
||||
eligible=clear_eligible,
|
||||
requires_confirmation=True,
|
||||
reason=(
|
||||
f"Binding classified as {binding_class.get('classification')}; clear is authorized."
|
||||
if clear_eligible
|
||||
else "Active worktree binding is clean, corroborated, or absent."
|
||||
),
|
||||
)
|
||||
)
|
||||
|
||||
# Playbook 2: Rebind Session Worktree
|
||||
rebind_eligible = bool(
|
||||
active_val
|
||||
or binding_class.get("classification") == stale_binding_recovery.CLASSIFICATION_UNVERIFIED_INHERITED
|
||||
)
|
||||
playbooks.append(
|
||||
PlaybookDescriptor(
|
||||
playbook_id=PLAYBOOK_REBIND_SESSION,
|
||||
label="Rebind Session Worktree",
|
||||
action_id=ACTION_REBIND_SESSION,
|
||||
description="Rebind or synchronize session worktree binding safely with active lease.",
|
||||
eligible=rebind_eligible,
|
||||
requires_confirmation=True,
|
||||
reason=(
|
||||
"Session worktree binding can be rebound to verified lease worktree."
|
||||
if rebind_eligible
|
||||
else "Session worktree is properly bound."
|
||||
),
|
||||
params_schema={"target_worktree": "string"},
|
||||
)
|
||||
)
|
||||
|
||||
# Playbook 3: Reconcile Cleanups
|
||||
reconcile_eligible = bool(hygiene.anomalies or any(e.classification in {"stale-clean", "detached-review"} for e in hygiene.entries))
|
||||
playbooks.append(
|
||||
PlaybookDescriptor(
|
||||
playbook_id=PLAYBOOK_RECONCILE_CLEANUPS,
|
||||
label="Trigger Reconciler Cleanups",
|
||||
action_id=ACTION_RECONCILE_CLEANUPS,
|
||||
description="Run sanctioned reconciler cleanup preview and apply for merged/superseded PR branches.",
|
||||
eligible=reconcile_eligible,
|
||||
requires_confirmation=True,
|
||||
reason=(
|
||||
f"Worktree hygiene scanner detected {len(hygiene.anomalies)} anomalies and cleanups needed."
|
||||
if reconcile_eligible
|
||||
else "No reconciler cleanups pending."
|
||||
),
|
||||
)
|
||||
)
|
||||
|
||||
# Playbook 4: Sanctioned Restart
|
||||
restart_eligible = bool(stale_runtime_obj.stale or contaminated)
|
||||
playbooks.append(
|
||||
PlaybookDescriptor(
|
||||
playbook_id=PLAYBOOK_SANCTIONED_RESTART,
|
||||
label="Sanctioned MCP Restart",
|
||||
action_id=sanctioned_restart.ACTION_RESTART_NAMESPACE,
|
||||
description="Restart MCP daemon via configured host supervisor without manual process kill.",
|
||||
eligible=restart_eligible,
|
||||
requires_confirmation=True,
|
||||
reason=(
|
||||
"Stale runtime or contamination detected; host supervisor restart available."
|
||||
if restart_eligible
|
||||
else "Runtime is healthy and clean."
|
||||
),
|
||||
params_schema={"namespace": "string", "mode": "restart|reload"},
|
||||
)
|
||||
)
|
||||
|
||||
clean = not reasons and not contaminated
|
||||
if contaminated:
|
||||
status = STATUS_BLOCKED_CONTAMINATION
|
||||
elif reasons:
|
||||
status = STATUS_ACTION_REQUIRED
|
||||
else:
|
||||
status = STATUS_HEALTHY
|
||||
|
||||
return RecoveryDiagnosis(
|
||||
status=status,
|
||||
clean=clean,
|
||||
stale_runtime=stale_runtime_dict,
|
||||
master_parity=parity_dict,
|
||||
stale_binding=binding_class,
|
||||
contamination=contamination_dict,
|
||||
worktree_anomalies=tuple(worktree_anomalies),
|
||||
playbooks=tuple(playbooks),
|
||||
reasons=tuple(reasons),
|
||||
)
|
||||
|
||||
|
||||
def build_recovery_preview(
|
||||
playbook_id: str,
|
||||
target: str | None = None,
|
||||
params: dict[str, Any] | None = None,
|
||||
principal: console_authz.Principal | None = None,
|
||||
env: dict[str, str] | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Generate dry-run preview & mutation ledger for a recovery playbook."""
|
||||
if playbook_id not in KNOWN_PLAYBOOKS:
|
||||
return {
|
||||
"allowed": False,
|
||||
"error": "unknown_playbook",
|
||||
"detail": f"Playbook {playbook_id!r} is not a registered recovery playbook.",
|
||||
}
|
||||
|
||||
action_id = PLAYBOOK_ACTIONS[playbook_id]
|
||||
action = console_authz.get_action(action_id)
|
||||
decision = console_authz.authorize(action_id, principal)
|
||||
# Preview and apply must answer the same question. ``execution_enabled`` was
|
||||
# a hardcoded False beside an authorization decision taken without
|
||||
# ``for_execution``, so the preview could not tell an operator *why*
|
||||
# execution was disabled — and the apply path did not ask at all.
|
||||
execution_decision = console_authz.authorize(
|
||||
action_id, principal, for_execution=True
|
||||
)
|
||||
phrase = confirmation_phrase(playbook_id, target)
|
||||
ledger = _build_ledger(playbook_id, target)
|
||||
|
||||
return {
|
||||
"playbook_id": playbook_id,
|
||||
"action_id": action_id,
|
||||
"target": target,
|
||||
"required_role": action.minimum_role if action else console_authz.OPERATOR,
|
||||
"required_permission": action.mcp_permission if action else "gitea.read",
|
||||
"requires_confirmation": True,
|
||||
"confirmation_phrase": phrase,
|
||||
"mutation_ledger": [asdict(entry) for entry in ledger],
|
||||
"authorization": decision.to_dict(),
|
||||
"execution_authorization": execution_decision.to_dict(),
|
||||
"params": dict(params or {}),
|
||||
"execution_enabled": bool(execution_decision.allowed),
|
||||
"execution_blocked_reason": (
|
||||
None if execution_decision.allowed else execution_decision.reason_code
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def execute_recovery_playbook(
|
||||
playbook_id: str,
|
||||
confirmation: str | None = None,
|
||||
target: str | None = None,
|
||||
params: dict[str, Any] | None = None,
|
||||
principal: console_authz.Principal | None = None,
|
||||
env: dict[str, str] | None = None,
|
||||
request_id: str | None = None,
|
||||
session_id: str | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Gated execution of a recovery playbook with audit logging and revalidation."""
|
||||
if playbook_id not in KNOWN_PLAYBOOKS:
|
||||
return {
|
||||
"success": False,
|
||||
"allowed": False,
|
||||
"error": "unknown_playbook",
|
||||
"detail": f"Playbook {playbook_id!r} is not known.",
|
||||
}
|
||||
|
||||
action_id = PLAYBOOK_ACTIONS[playbook_id]
|
||||
# The mapping the playbooks actually mutate. ``dict(os.environ)`` produced a
|
||||
# throwaway copy: every env playbook wrote to it, verified against it, and
|
||||
# left the running daemon bound to the value it claimed to have fixed.
|
||||
mutation_env: Any = env if env is not None else os.environ
|
||||
|
||||
# 1. Authorization check — ``for_execution=True`` is what arms the phase
|
||||
# gate (console_authz.authorize only applies it in that branch). Without it
|
||||
# a phase-2 write executed while the console was in phase 1.
|
||||
decision = console_authz.authorize(action_id, principal, for_execution=True)
|
||||
if not decision.allowed:
|
||||
console_audit.record_event(
|
||||
action_id=action_id,
|
||||
result=console_audit.RESULT_DENIED,
|
||||
principal=principal,
|
||||
target={"playbook_id": playbook_id, "target": target},
|
||||
reason_code=decision.reason_code,
|
||||
detail=decision.detail,
|
||||
request_id=request_id,
|
||||
session_id=session_id,
|
||||
)
|
||||
return {
|
||||
"success": False,
|
||||
"allowed": False,
|
||||
"error": decision.reason_code,
|
||||
"detail": decision.detail,
|
||||
}
|
||||
|
||||
# 2. Confirmation phrase check
|
||||
if not confirmation_matches(playbook_id, confirmation, target):
|
||||
expected = confirmation_phrase(playbook_id, target)
|
||||
detail = f"Confirmation phrase mismatch. Expected: {expected!r}"
|
||||
console_audit.record_event(
|
||||
action_id=action_id,
|
||||
result=console_audit.RESULT_DENIED,
|
||||
principal=principal,
|
||||
target={"playbook_id": playbook_id, "target": target},
|
||||
reason_code="confirmation_mismatch",
|
||||
detail=detail,
|
||||
request_id=request_id,
|
||||
session_id=session_id,
|
||||
)
|
||||
return {
|
||||
"success": False,
|
||||
"allowed": False,
|
||||
"error": "confirmation_mismatch",
|
||||
"detail": detail,
|
||||
"expected_confirmation_phrase": expected,
|
||||
}
|
||||
|
||||
# 3. Contamination rule (#630) check
|
||||
role_str = principal.role if principal else None
|
||||
contamination_marker = load_active_contamination_marker(env=env)
|
||||
contam = runtime_recovery_guard.assess_contamination_gate(
|
||||
contamination_marker,
|
||||
task=CONTAMINATION_GATED_TASK,
|
||||
actual_role=role_str,
|
||||
)
|
||||
if contam.get("block"):
|
||||
if playbook_id != PLAYBOOK_RECONCILE_CLEANUPS:
|
||||
detail = "Runtime is contaminated by a manual process kill (#630). Run reconciler cleanup playbook first."
|
||||
console_audit.record_event(
|
||||
action_id=action_id,
|
||||
result=console_audit.RESULT_DENIED,
|
||||
principal=principal,
|
||||
target={"playbook_id": playbook_id, "target": target},
|
||||
reason_code="contaminated_runtime",
|
||||
detail=detail,
|
||||
request_id=request_id,
|
||||
session_id=session_id,
|
||||
)
|
||||
return {
|
||||
"success": False,
|
||||
"allowed": False,
|
||||
"error": "contaminated_runtime",
|
||||
"detail": detail,
|
||||
}
|
||||
|
||||
# 4. Execute playbook action
|
||||
applied_result: dict[str, Any] = {"performed": False}
|
||||
if playbook_id == PLAYBOOK_CLEAR_STALE_BINDING:
|
||||
binding_before = _active_binding(mutation_env)
|
||||
diagnosis = diagnose_recovery(env=env)
|
||||
plan = stale_binding_recovery.plan_recovery(diagnosis.stale_binding)
|
||||
applied_result = stale_binding_recovery.apply_recovery(plan, env=mutation_env)
|
||||
binding_after = _active_binding(mutation_env)
|
||||
applied_result = {
|
||||
**applied_result,
|
||||
"binding_before": binding_before,
|
||||
"binding_after": binding_after,
|
||||
"binding_changed": binding_before != binding_after,
|
||||
}
|
||||
# A clear that did not clear is not a success, whatever the plan said.
|
||||
if not applied_result["binding_changed"]:
|
||||
applied_result["performed"] = False
|
||||
applied_result.setdefault("reasons", []).append(
|
||||
"clear_stale_binding did not change the live worktree binding"
|
||||
)
|
||||
elif playbook_id == PLAYBOOK_REBIND_SESSION:
|
||||
target_wt = target or (params or {}).get("target_worktree")
|
||||
if target_wt:
|
||||
binding_before = _active_binding(mutation_env)
|
||||
mutation_env[stale_binding_recovery.ACTIVE_WORKTREE_ENV] = target_wt
|
||||
binding_after = _active_binding(mutation_env)
|
||||
applied_result = {
|
||||
"performed": binding_after == target_wt,
|
||||
"rebound_worktree": target_wt,
|
||||
"binding_before": binding_before,
|
||||
"binding_after": binding_after,
|
||||
"binding_changed": binding_before != binding_after,
|
||||
"cleared_stale": binding_before != binding_after,
|
||||
}
|
||||
if binding_after != target_wt:
|
||||
applied_result["reasons"] = [
|
||||
"rebind_session_worktree did not take effect on the live "
|
||||
"environment"
|
||||
]
|
||||
else:
|
||||
applied_result = {
|
||||
"performed": False,
|
||||
"reason": "No target_worktree specified for rebind.",
|
||||
}
|
||||
elif playbook_id == PLAYBOOK_RECONCILE_CLEANUPS:
|
||||
# ``merged_cleanup_reconcile`` exposes the building blocks only; the
|
||||
# orchestrator is the MCP tool. The previous call named a function that
|
||||
# does not exist, and a bare ``except`` turned the AttributeError into a
|
||||
# generic failure, so this playbook could never succeed. Imported lazily
|
||||
# because the MCP server module is large and binds FastMCP at import.
|
||||
try:
|
||||
import gitea_mcp_server
|
||||
|
||||
snapshot = gitea_mcp_server.gitea_reconcile_merged_cleanups(
|
||||
dry_run=False,
|
||||
execute_confirmed=True,
|
||||
remote=(params or {}).get("remote") or _console_remote(env),
|
||||
org=(params or {}).get("org"),
|
||||
repo=(params or {}).get("repo"),
|
||||
)
|
||||
performed_reconcile = bool(snapshot.get("success"))
|
||||
applied_result = {
|
||||
"performed": performed_reconcile,
|
||||
"reconciled_count": len(snapshot.get("entries") or []),
|
||||
"snapshot": snapshot,
|
||||
}
|
||||
if not performed_reconcile:
|
||||
applied_result["reasons"] = list(snapshot.get("reasons") or [])
|
||||
except Exception as exc: # noqa: BLE001 — surfaced with its type
|
||||
applied_result = {
|
||||
"performed": False,
|
||||
"error": str(exc),
|
||||
"error_type": type(exc).__name__,
|
||||
}
|
||||
elif playbook_id == PLAYBOOK_SANCTIONED_RESTART:
|
||||
ns = target or (params or {}).get("namespace", "gitea-author")
|
||||
md = (params or {}).get("mode", sanctioned_restart.MODE_RESTART)
|
||||
restart_res = sanctioned_restart.execute_restart(
|
||||
namespace=ns,
|
||||
mode=md,
|
||||
principal=principal,
|
||||
confirmation=f"{md} {ns}",
|
||||
# Without the marker the stricter guard at sanctioned_restart.py:375
|
||||
# never fires and a restart can launder a contaminated runtime.
|
||||
contamination_marker=contamination_marker,
|
||||
env=mutation_env,
|
||||
request_id=request_id,
|
||||
session_id=session_id,
|
||||
)
|
||||
applied_result = restart_res
|
||||
|
||||
# ``allowed`` is not ``performed``: execute_restart documents that success is
|
||||
# False in both directions because the host supervisor still has to act.
|
||||
performed = bool(applied_result.get("performed"))
|
||||
|
||||
# 5. Record Audit Log
|
||||
audit_record = console_audit.record_event(
|
||||
action_id=action_id,
|
||||
result=console_audit.RESULT_ALLOWED if performed else console_audit.RESULT_DENIED,
|
||||
principal=principal,
|
||||
target={"playbook_id": playbook_id, "target": target},
|
||||
reason_code="recovery_executed" if performed else "recovery_failed",
|
||||
detail=f"Executed recovery playbook {playbook_id}",
|
||||
request_id=request_id,
|
||||
session_id=session_id,
|
||||
metadata={"applied_result": applied_result},
|
||||
)
|
||||
|
||||
# 6. Post-recovery verification recheck.
|
||||
#
|
||||
# Re-read state rather than re-reading the mapping the mutation just wrote:
|
||||
# verifying the mutated copy confirmed changes that never reached the
|
||||
# process. Passing ``env`` through means a caller-supplied mapping is the
|
||||
# live one for that caller, and ``None`` re-reads ``os.environ`` fresh.
|
||||
post_verification = verify_post_recovery(env=env)
|
||||
|
||||
return {
|
||||
"success": performed,
|
||||
"allowed": True,
|
||||
"playbook_id": playbook_id,
|
||||
"action_id": action_id,
|
||||
"applied_result": applied_result,
|
||||
"audit": audit_record,
|
||||
"post_recovery_verification": post_verification,
|
||||
}
|
||||
|
||||
|
||||
def verify_post_recovery(
|
||||
repo_path: Path | str | None = None, env: dict[str, str] | None = None
|
||||
) -> dict[str, Any]:
|
||||
"""Revalidate control-plane state post-recovery before clean status."""
|
||||
diag = diagnose_recovery(repo_path, env)
|
||||
classification = diag.stale_binding.get("classification")
|
||||
# ``not clear_eligible`` also reads clean for every binding recovery is not
|
||||
# allowed to touch — an unverified inherited binding is unproven, not clean.
|
||||
binding_clean = (
|
||||
not diag.stale_binding.get("clear_eligible")
|
||||
and classification != stale_binding_recovery.CLASSIFICATION_UNVERIFIED_INHERITED
|
||||
)
|
||||
return {
|
||||
"clean": diag.clean,
|
||||
"status": diag.status,
|
||||
"stale_runtime_clean": not diag.stale_runtime.get("stale"),
|
||||
"binding_clean": binding_clean,
|
||||
"binding_classification": classification,
|
||||
# The gate returns ``block``; it has never returned ``contaminated``, so
|
||||
# reading that key reported every runtime clean unconditionally.
|
||||
"contamination_clean": not diag.contamination.get("block"),
|
||||
"anomalies_count": len(diag.worktree_anomalies),
|
||||
"reasons": list(diag.reasons),
|
||||
}
|
||||
@@ -178,13 +178,6 @@ def build_action_registry() -> ActionRegistry:
|
||||
("system.restart_namespace", "Restart MCP namespace",
|
||||
"restart_namespace", "host.supervisor_restart",
|
||||
"Restart one MCP namespace via the host supervisor."),
|
||||
# #644: Phase 2 recovery playbooks & controls.
|
||||
("system.clear_stale_binding", "Clear stale binding", "clear_stale_binding",
|
||||
"console.clear_stale_binding", "Clear provably stale or superseded env binding."),
|
||||
("system.rebind_session_worktree", "Rebind session worktree", "rebind_session_worktree",
|
||||
"console.rebind_session_worktree", "Rebind session worktree to verified lease."),
|
||||
("system.reconcile_cleanups", "Reconcile cleanups", "reconcile_cleanups",
|
||||
"console.reconcile_cleanups", "Run reconciler cleanup for merged or superseded PRs."),
|
||||
)
|
||||
actions = tuple(
|
||||
GatedAction(
|
||||
|
||||
@@ -0,0 +1,713 @@
|
||||
"""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()
|
||||
@@ -0,0 +1,196 @@
|
||||
"""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)
|
||||
+6
-12
@@ -4,12 +4,11 @@ Single source of truth for the console navigation so ``webui/layout.py`` and
|
||||
the ``webui/app.py`` route table stay aligned with epic #631. Read-only: every
|
||||
destination is a GET view or a Phase 1 placeholder. No mutation links.
|
||||
|
||||
Nav groups follow the #631 Phase 1 information architecture: Health, Traffic,
|
||||
Nav groups follow the #631 information architecture: Health, Traffic,
|
||||
Runtime/Sessions, Projects, Inventory, Timeline, Policy (placeholder), and
|
||||
Insights (placeholder), joined by the Phase 3 Gitea linkage group (#645).
|
||||
Later-phase surfaces are declared as ``stub`` items and
|
||||
backed by ``STUB_PAGES`` so their nav links resolve to a graceful placeholder
|
||||
instead of a 404.
|
||||
Gitea linkage (#645) plus Phase 4 Insights/Providers (#650). Later-phase
|
||||
surfaces are declared as ``stub`` items and backed by ``STUB_PAGES`` so their
|
||||
nav links resolve to a graceful placeholder instead of a 404.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -46,7 +45,6 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
|
||||
NavItem("/queue", "Queue"),
|
||||
NavItem("/leases", "Leases"),
|
||||
NavItem("/actions", "Actions"),
|
||||
NavItem("/requests", "Requests"),
|
||||
)),
|
||||
NavGroup("Runtime/Sessions", (
|
||||
NavItem("/runtime", "Runtime health"),
|
||||
@@ -70,7 +68,8 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
|
||||
NavItem("/prompts", "Prompts"),
|
||||
)),
|
||||
NavGroup("Insights", (
|
||||
NavItem("/insights", "Insights", "stub"),
|
||||
NavItem("/insights", "Insights"),
|
||||
NavItem("/providers", "Providers"),
|
||||
NavItem("/analytics", "Analytics"),
|
||||
NavItem("/audit", "Audit"),
|
||||
)),
|
||||
@@ -94,11 +93,6 @@ STUB_PAGES: dict[str, tuple[str, str]] = {
|
||||
"Policy",
|
||||
"Capability and role policy surface. Placeholder until a later phase.",
|
||||
),
|
||||
"/insights": (
|
||||
"Insights",
|
||||
"Aggregate operational insights and trends. Placeholder until a later "
|
||||
"phase.",
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -1,164 +0,0 @@
|
||||
"""HTML views for the operator request surface (#643).
|
||||
|
||||
The form is deliberately a *preview* form. It has no initiate button, because
|
||||
initiating requires a confirmed POST to ``/api/v1/requests/apply`` and a stray
|
||||
form submission must not be able to produce one by accident.
|
||||
|
||||
Nothing rendered here is trusted input: every interpolated value is escaped,
|
||||
and the page renders only values the service already produced rather than
|
||||
echoing a raw request body back.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import html
|
||||
import json
|
||||
from typing import Any
|
||||
|
||||
from webui.layout import render_page
|
||||
from webui.request_service import (
|
||||
REQUESTABLE_ROLES,
|
||||
WORK_KINDS,
|
||||
RequestError,
|
||||
RequestPreview,
|
||||
)
|
||||
|
||||
REQUESTS_PATH = "/requests"
|
||||
PREVIEW_API_PATH = "/api/v1/requests/preview"
|
||||
APPLY_API_PATH = "/api/v1/requests/apply"
|
||||
|
||||
|
||||
def _escape(text: Any) -> str:
|
||||
return html.escape(str(text if text is not None else ""), quote=True)
|
||||
|
||||
|
||||
REQUEST_PAGE_STYLES = """
|
||||
<style>
|
||||
.request-form { display: grid; gap: 0.75rem; max-width: 44rem; }
|
||||
.request-form label { display: grid; gap: 0.25rem; font-size: 0.9rem; }
|
||||
.request-check { margin: 0.35rem 0; }
|
||||
.request-check .verdict-ok { color: var(--accent); }
|
||||
.request-check .verdict-fail { color: #d14; }
|
||||
.request-prohibited code { margin-right: 0.4rem; }
|
||||
</style>
|
||||
"""
|
||||
|
||||
|
||||
def _options(values: tuple[str, ...], selected: Any) -> str:
|
||||
return "".join(
|
||||
f"<option value='{_escape(value)}'"
|
||||
+ (" selected" if selected == value else "")
|
||||
+ f">{_escape(value)}</option>"
|
||||
for value in values
|
||||
)
|
||||
|
||||
|
||||
def _form(values: dict[str, Any] | None = None) -> str:
|
||||
current = dict(values or {})
|
||||
number = current.get("work_number")
|
||||
return (
|
||||
f"<form class='request-form' method='post' action='{REQUESTS_PATH}'>"
|
||||
"<label>Desired role<select name='desired_role'>"
|
||||
f"{_options(REQUESTABLE_ROLES, current.get('desired_role'))}"
|
||||
"</select></label>"
|
||||
"<label>Work kind<select name='work_kind'>"
|
||||
f"{_options(WORK_KINDS, current.get('work_kind'))}"
|
||||
"</select></label>"
|
||||
"<label>Issue or PR number"
|
||||
"<input type='number' name='work_number' min='1' required "
|
||||
f"value='{_escape(number) if number else ''}'></label>"
|
||||
"<label>Intent summary"
|
||||
"<input type='text' name='intent_summary' maxlength='500' required "
|
||||
f"value='{_escape(current.get('intent_summary'))}'></label>"
|
||||
"<label>Expected head SHA <span class='muted'>(PR work only)</span>"
|
||||
"<input type='text' name='expected_head_sha' "
|
||||
f"value='{_escape(current.get('expected_head_sha'))}'></label>"
|
||||
"<button type='submit' class='copy-btn'>Preview request</button>"
|
||||
"<p class='muted meta'>Preview is read-only and creates no assignment. "
|
||||
f"Initiating requires a confirmed POST to <code>{APPLY_API_PATH}</code>."
|
||||
"</p>"
|
||||
"</form>"
|
||||
)
|
||||
|
||||
|
||||
def _checks_block(preview: RequestPreview) -> str:
|
||||
rows = []
|
||||
for check in preview.checks:
|
||||
verdict = "PASS" if check.ok else "FAIL"
|
||||
css = "verdict-ok" if check.ok else "verdict-fail"
|
||||
rows.append(
|
||||
"<li class='request-check'>"
|
||||
f"<span class='{css}'><strong>{verdict}</strong></span> "
|
||||
f"<code>{_escape(check.name)}</code> — {_escape(check.detail)} "
|
||||
f"<span class='muted meta'>({_escape(check.reason_code)})</span>"
|
||||
"</li>"
|
||||
)
|
||||
return "<ul>" + "".join(rows) + "</ul>"
|
||||
|
||||
|
||||
def _preview_block(preview: RequestPreview) -> str:
|
||||
verdict = "AUTHORIZED" if preview.authorized else "DENIED"
|
||||
prohibited = "".join(
|
||||
f"<code>{_escape(action)}</code>" for action in preview.prohibited_actions
|
||||
)
|
||||
request = preview.request
|
||||
evidence = json.dumps(preview.allocator_evidence, indent=2, default=str)
|
||||
return (
|
||||
"<h3>Intent preview</h3>"
|
||||
f"<p><strong>{verdict}</strong> — {_escape(preview.detail)}</p>"
|
||||
"<p class='meta'>"
|
||||
f"Role <code>{_escape(request.desired_role)}</code> · "
|
||||
f"{_escape(request.work_kind)} <code>{_escape(request.display_ref)}</code>"
|
||||
f" · profile <code>{_escape(preview.required_profile)}</code> · "
|
||||
f"namespace <code>{_escape(preview.required_namespace)}</code> · "
|
||||
f"permission <code>{_escape(preview.required_permission)}</code>"
|
||||
"</p>"
|
||||
f"<p>Intent: {_escape(request.intent_summary)}</p>"
|
||||
f"{_checks_block(preview)}"
|
||||
f"<p><strong>Next safe action:</strong> "
|
||||
f"{_escape(preview.next_safe_action)}</p>"
|
||||
"<p class='request-prohibited'><strong>Prohibited for this role:</strong> "
|
||||
+ (prohibited or "<span class='muted'>none declared</span>")
|
||||
+ "</p>"
|
||||
"<p class='muted meta'>Correlation id "
|
||||
f"<code>{_escape(preview.correlation_id)}</code></p>"
|
||||
"<details><summary>Allocator evidence</summary>"
|
||||
f"<pre class='prompt-text'>{_escape(evidence)}</pre>"
|
||||
"</details>"
|
||||
)
|
||||
|
||||
|
||||
def _error_block(error: RequestError) -> str:
|
||||
field = (
|
||||
f"<p class='meta'>Field: <code>{_escape(error.field_name)}</code></p>"
|
||||
if error.field_name
|
||||
else ""
|
||||
)
|
||||
return (
|
||||
"<h3>Request rejected</h3>"
|
||||
f"<p><strong>{_escape(error.reason_code)}</strong> — "
|
||||
f"{_escape(error.detail)}</p>{field}"
|
||||
)
|
||||
|
||||
|
||||
def render_requests_page(
|
||||
*,
|
||||
preview: RequestPreview | None = None,
|
||||
error: RequestError | None = None,
|
||||
submitted: dict[str, Any] | None = None,
|
||||
) -> str:
|
||||
"""Render the request form, plus a preview or rejection when one exists."""
|
||||
body = (
|
||||
"<h2>Requests</h2>"
|
||||
"<p>Submit a work request — desired role, issue or PR, and intent — "
|
||||
"and see whether it would be authorized before anything is reserved. "
|
||||
"Initiation goes through the allocator (#600/#613); this console never "
|
||||
"self-selects work, never approves, and never merges.</p>"
|
||||
+ _form(submitted)
|
||||
+ (_error_block(error) if error is not None else "")
|
||||
+ (_preview_block(preview) if preview is not None else "")
|
||||
+ f"<p class='meta'><a href='{PREVIEW_API_PATH}'>Preview API</a> · "
|
||||
"<a href='/api/console/security-model'>RBAC model</a></p>"
|
||||
+ REQUEST_PAGE_STYLES
|
||||
)
|
||||
return render_page(title="Requests", body_html=body)
|
||||
@@ -270,53 +270,26 @@ def _probe_error_card(snapshot: SystemHealthSnapshot) -> str:
|
||||
|
||||
|
||||
def _recovery_card() -> str:
|
||||
"""Sanctioned recovery controls & playbooks (#644, Phase 2)."""
|
||||
try:
|
||||
from webui import console_recovery
|
||||
diag = console_recovery.diagnose_recovery()
|
||||
# Every other card in this file escapes at the interpolation boundary.
|
||||
# This one did not, and it is where a #630 marker's operator-supplied
|
||||
# command_summary lands once the contamination gate is wired.
|
||||
status_badge = (
|
||||
f"<span class='status-pill {_esc(diag.status)}'>{_esc(diag.status)}</span>"
|
||||
)
|
||||
playbook_lis = ""
|
||||
for pb in diag.playbooks:
|
||||
elig = "eligible" if pb.eligible else "disabled"
|
||||
playbook_lis += (
|
||||
f"<li><strong>{_esc(pb.label)}</strong> "
|
||||
f"(<code>{_esc(pb.playbook_id)}</code>) — "
|
||||
f"<span class='badge {elig}'>{elig}</span>: {_esc(pb.description)} "
|
||||
f"<em class='muted'>({_esc(pb.reason)})</em></li>"
|
||||
)
|
||||
reasons_html = ""
|
||||
if diag.reasons:
|
||||
items = "".join(f"<li>{_esc(r)}</li>" for r in diag.reasons)
|
||||
reasons_html = f"<ul class='reasons'>{items}</ul>"
|
||||
else:
|
||||
reasons_html = "<p class='clean-note'>No recovery actions currently required. Control plane is healthy.</p>"
|
||||
|
||||
return (
|
||||
"<section class='health-card recovery-card'>"
|
||||
f"<h3>Sanctioned Recovery Controls (Phase 2 #644) {status_badge}</h3>"
|
||||
"<p class='muted'>Guided recovery wizard: Diagnose → Preview → Confirm → Verify. "
|
||||
"Reconnect the MCP client from the IDE, then re-run the blocked cycle. "
|
||||
"Never kill the daemon process manually: unmanaged kills are recorded as runtime contamination (#630).</p>"
|
||||
f"{reasons_html}"
|
||||
"<h4>Available Recovery Playbooks</h4>"
|
||||
f"<ul class='playbooks-list'>{playbook_lis}</ul>"
|
||||
"<p class='meta'>APIs: <code>/api/v1/system/recovery/diagnose</code>, "
|
||||
"<code>/api/v1/system/recovery/preview</code>, <code>/api/v1/system/recovery/apply</code>, "
|
||||
"<code>/api/v1/system/recovery/verify</code>.</p>"
|
||||
"</section>"
|
||||
)
|
||||
except Exception as exc:
|
||||
return (
|
||||
"<section class='health-card'>"
|
||||
"<h3>Sanctioned Recovery Controls (Phase 2 #644)</h3>"
|
||||
f"<p class='error'>Recovery diagnostics unavailable: {_esc(exc)}</p>"
|
||||
"</section>"
|
||||
)
|
||||
"""Sanctioned recovery pointers only — never a manual process kill (#630)."""
|
||||
return (
|
||||
"<section class='health-card'>"
|
||||
"<h3>Recovery</h3>"
|
||||
"<p class='muted'>This dashboard is read-only. Restart and reload "
|
||||
"controls arrive in Phase 2 (#642); until then recovery runs through "
|
||||
"the sanctioned client reconnect / operator restart path.</p>"
|
||||
"<ul class='reasons'>"
|
||||
"<li><a href='/runtime'>Runtime health</a> — active profile, workflow "
|
||||
"hashes, and shell health.</li>"
|
||||
"<li><a href='/sessions'>Runtime and sessions</a> — namespaces, session "
|
||||
"rows, worktree bindings, and contamination markers (#641).</li>"
|
||||
"<li>Reconnect the MCP client from the IDE, then re-run the blocked "
|
||||
"cycle. Never kill the daemon process manually: unmanaged kills are "
|
||||
"recorded as runtime contamination (#630).</li>"
|
||||
"<li>See <code>docs/webui-local-dev.md</code> for the documented "
|
||||
"recovery sequence.</li>"
|
||||
"</ul>"
|
||||
"</section>"
|
||||
)
|
||||
|
||||
|
||||
def render_system_health_page(snapshot: SystemHealthSnapshot) -> str:
|
||||
|
||||
@@ -201,16 +201,6 @@ def _candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate
|
||||
return candidates
|
||||
|
||||
|
||||
def candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate]:
|
||||
"""Public alias for :func:`_candidates_from_queue_snapshot` (#643).
|
||||
|
||||
The request-initiation service ranks the same candidate set this view
|
||||
renders, so both must agree on how a queue row becomes a candidate. One
|
||||
construction, two callers — not two that can drift apart.
|
||||
"""
|
||||
return _candidates_from_queue_snapshot(q_snap)
|
||||
|
||||
|
||||
def _claim_lease_records(inventory: dict[str, Any] | None) -> list[dict[str, Any]]:
|
||||
"""Normalize ``build_claim_inventory`` entries into lease records.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user