Compare commits

..
Author SHA1 Message Date
sysadmin 0f9390aab4 Merge remote-tracking branch 'prgs/master' into feat/issue-666-concurrent-mcp-restart-tests 2026-07-25 18:43:28 -04:00
sysadmin d7ad2838ec Merge pull request 'docs(incident): retroactive audit for direct-to-master commit 2fa97c26 (#670)' (#915) from fix/issue-670-direct-master-incident into master 2026-07-25 17:40:23 -05:00
sysadmin c6d68dbc7b Merge pull request 'feat(webui): notifications and human-attention routing (#648)' (#905) from feat/issue-648-notifications-console into master 2026-07-25 17:40:03 -05:00
sysadmin c83a10d7c2 Merge remote-tracking branch 'prgs/master' into feat/issue-648-notifications-console 2026-07-25 18:37:17 -04:00
sysadmin ca22c326a4 Merge pull request 'docs(architecture): ADR for high-availability and rolling-restart MCP architecture (Closes #668)' (#912) from docs/issue-668-mcp-ha-rolling-restart into master 2026-07-25 17:34:31 -05:00
sysadmin 3bbe6df6c7 Merge pull request 'feat(webui): add read-only console restart status and impact controls (Closes #667)' (#911) from feat/issue-667-console-restart-controls into master 2026-07-25 17:34:12 -05:00
jcwalker3 71031c812e Merge branch 'master' into feat/issue-648-notifications-console 2026-07-25 17:29:57 -05:00
jcwalker3 e43ddd3cbe docs(incident): retroactive audit for direct-to-master commit 2fa97c26 (#670) 2026-07-25 17:29:34 -05:00
sysadmin a64ba08e27 fix(webui): address #905 REQUEST_CHANGES on notifications classifier
B1: classify_attention_event uses structured flags/category only — never
substring-match human-authored title/summary for escalation.

B2: make notification ids unique across probe_errors and collisions
(include loop index / kind).

B3: do not assign probe_errors to fetch_error (avoids false Fetch Warning
and double-reporting).

Regression tests cover all three blockers.

Refs #648
2026-07-25 18:26:46 -04:00
sysadmin 26f54851d1 Merge pull request 'feat(mcp): enforce scoped recovery playbook before full restart (Closes #669)' (#914) from feat/issue-669-scoped-component-recovery into master 2026-07-25 17:14:21 -05:00
sysadmin f02a2dc030 Merge pull request 'feat(webui): Web Console stale-runtime recovery & reconciliation controls (Phase 2) (#644)' (#903) from feat/issue-644-console-recovery into master 2026-07-25 17:14:06 -05:00
sysadminandClaude Opus 4.8 bb8c3a537b merge(master): resolve PR #905 conflicts with requests/linkage
Keep notifications (#648) routes and nav alongside master requests (#643)
and other base updates.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 18:11:06 -04:00
sysadminandClaude Opus 4.8 04ae3532cc merge(master): resolve #903 conflicts with #643 request initiation
Keep #644 recovery console actions and #643 initiate_workflow side by side
in console_authz and authz audit docs after merging latest master.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 18:03:04 -04:00
sysadmin 7bb5ff4719 Merge pull request 'feat(webui): request preview, authorization, and workflow initiation (Closes #643)' (#902) from feat/issue-643-request-preview-initiate into master 2026-07-25 17:02:17 -05:00
sysadmin 5b7ceefa9a Merge remote-tracking branch 'prgs/master' into feat/issue-644-console-recovery 2026-07-25 18:00:31 -04:00
sysadminandClaude Opus 4.8 4a2fae8495 fix(webui): arm the recovery gates and make the playbooks reach the process (#644)
Reviewer REQUEST_CHANGES on PR #903 at head 1c88b87 raised five blockers, all
reproduced by executing that head. The shared shape: a write path that declared
itself gated, audited, and verified, but never armed the gate, mutated a copy of
the state it claimed to fix, and then verified against that same copy.

B1 - the apply path never asked the execution gate.
execute_recovery_playbook called console_authz.authorize with the default
for_execution=False, and the phase branch only fires when it is True. ACTIVE_PHASE
is 1 and every new action is phase 2, so an operator executed a phase-2 write
through POST /api/v1/system/recovery/apply while build_recovery_preview reported
execution_enabled false. The call now passes for_execution=True and surfaces the
phase_not_active refusal. Preview reports the same decision under
execution_authorization / execution_blocked_reason instead of a hardcoded False it
could not explain.

B2 - both env playbooks mutated a discarded copy and verified against it.
source_env = dict(os.environ) meant clear_stale_binding and rebind_session_worktree
never touched the running process, and verify_post_recovery(env=source_env)
re-diagnosed the same copy, confirming a change that had not happened. Mutations
now target the live mapping (apply_recovery's sanctioned env=None -> os.environ
path, #702 AC2) and verification re-reads state rather than the mutated input.
binding_before / binding_after / binding_changed are returned, and a playbook that
changed nothing reports performed: false. verify_post_recovery no longer reads an
unverified_inherited binding as clean, because unproven is not clean.

B3 - the reconcile playbook called a function that does not exist.
merged_cleanup_reconcile.reconcile_merged_cleanups is absent from that module and a
bare except turned the AttributeError into a generic failure, so the playbook could
never succeed. It now calls gitea_mcp_server.gitea_reconcile_merged_cleanups, the
real orchestrator, imported lazily; failures carry error_type. task_capability_map
declared gitea.pr.close for reconcile_cleanups while the entry point gates on
gitea.read; the two authority statements are reconciled to the one that is enforced.

B4 - the #630 contamination integration could not block.
assess_contamination_gate was fed marker=None, which short-circuits to block: False
on its first statement; the task passed was a console action id outside
CONTAMINATION_GATED_TASKS; and the result was read through a "contaminated" key the
gate never returns, making STATUS_BLOCKED_CONTAMINATION unreachable. The live marker
now comes from the #641 session inventory reader, the gated task key
console_recovery_apply is added to CONTAMINATION_GATED_TASKS, every read uses the
"block" key the gate actually returns, and the marker is forwarded to
sanctioned_restart.execute_restart so a restart cannot launder a contaminated
runtime. The reconciler cleanup playbook stays exempt as the designated remedy.

B5 - the parity baseline was captured from the head it was compared against.
capture_startup_parity(root, head=checkout_head) stores the head verbatim, so
in_parity was structurally incapable of being false, and live_remote_head was never
passed. The baseline is now the daemon start head that assess_stale_runtime already
returns, and the #610 live-remote dimension is restored.

Also: _recovery_card was the one renderer in system_health_views.py interpolating
without _esc(), and it is where a marker's operator-supplied command_summary lands
once B4 is wired; it now escapes, including the except branch. Docs no longer claim
apply enforces master parity or that verify asserts clean: true, and the absolute
file:///Users/... links are relative.

Tests: the two that asserted the defects as intended are inverted -
test_api_recovery_apply_with_dev_auth asserted the phase-gate bypass, and the rebind
test asserted the input echoed back. Added coverage per blocker, including a no-op
detection test that fails when a playbook reports success without changing anything,
the previously untested reconcile playbook, contamination block and remedy-exemption
tests, and parity baseline/live-remote tests. Both new guards were mutation-verified:
disarming for_execution fails 2 tests, restoring the env copy fails 2 tests.

Validation: WEBUI_TEST_OFFLINE=1 ../../venv/bin/python -m pytest tests/ -q from
branches/feat-issue-644 gives 27 failed / 5291 passed / 6 skipped / 953 subtests;
the same command from branches/baseline-master-76f293e at 76f293eb28 gives
28 failed / 5262 passed / 6 skipped / 926 subtests. Suites run one at a time.
comm of the sorted FAILED lines shows no new signature at the head. The single
absent signature, test_workspace_guard_alignment.py::
TestRuntimeContextGuardAlignment::test_declared_branches_worktree_passes_when_mcp_root_differs,
is suite-order dependent: that file passes 9/9 in isolation at both revisions.

Closes #644

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 17:42:49 -04:00
sysadminandGrok 4.5 461e1dac78 feat(mcp): enforce scoped recovery playbook before full restart (Closes #669)
Add recovery_playbook.py with the narrow-to-broad recovery ladder, symptom
routing, attempt-log helpers, and escalation metrics. Wire the attempt-log
gate into restart_coordinator so rolling/full/host restarts require prior
insufficient narrower attempts (or break-glass). Document the ladder and
update gitea_request_mcp_restart for prior_recovery_attempts_json.

Co-Authored-By: Grok 4.5 <[email protected]>
2026-07-25 17:35:11 -04:00
sysadmin 6010f4295b docs(architecture): add ADR for HA and rolling-restart MCP architecture (Closes #668) 2026-07-25 17:27:50 -04:00
sysadminandClaude Opus 4.8 9b8e315b49 docs(webui): document the read-only restart console surface (#667)
Records the two GET routes, what each panel consumes, and the three
properties the surface is held to: an unreadable source reports
unavailable rather than green, authorization is probed with
for_execution=True so a Phase 1 refusal is never shown as an allow, and
the control-plane database is opened mode=ro so reading status never
creates it.

Refs #655 #642 #658 #661 #662 #663 #633 #652 #653 #664

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 17:26:30 -04:00
sysadmin 9a01543477 feat(webui): add read-only console restart status and impact controls (Closes #667) 2026-07-25 17:23:52 -04:00
sysadmin 59aab06fe1 feat(tests): add concurrent-session MCP restart safety tests (Closes #666) 2026-07-25 17:14:14 -04:00
sysadmin 4f06d30e07 feat(webui): implement notifications and human-attention routing (#648) 2026-07-25 16:41:16 -04:00
sysadmin 1c88b87ec5 feat(webui): implement Phase 2 recovery controls & playbooks (#644) 2026-07-25 08:19:06 -04:00
jcwalker3 b993ad1c64 Merge branch 'master' into feat/issue-643-request-preview-initiate 2026-07-25 07:13:21 -05:00
jcwalker3 d0006e9f71 Merge branch 'master' into feat/issue-643-request-preview-initiate 2026-07-25 06:42:41 -05:00
jcwalker3 6da68fffb8 Merge branch 'master' into feat/issue-643-request-preview-initiate 2026-07-25 02:44:58 -05:00
sysadminandClaude Opus 5 53ce1b1a5e fix(webui): compensate stray allocations and keep preview side-effect free
Addresses both blockers from the PR #902 review (review 589) for issue #643.

B1 — an allocator-created assignment could be orphaned and reported as no
mutation.

apply_request re-previews, then re-runs the allocator with apply=True. The CAS
fingerprint hashes only {kind, number} plus exclusions, so a competing lease
taken on the requested unit inside the window leaves the fingerprint identical:
the pin passes, the selection loop skips the now-claimed unit and commits an
assignment on the *next* one, and _selection_matches then fails on egress. The
old code returned mutation_performed False with that lease still committed and
owned by a synthetic session nothing heartbeats. The window contains a second
full load_queue_snapshot(), so it is seconds wide, and foreign sessions acting
on this repo concurrently are an observed condition.

The egress mismatch now releases the assignment the allocator created before
refusing. When the release succeeds the refusal reports mutation_performed
False and a compensation record; when it fails the response carries
mutation_performed True, an explicit orphaned_assignment, and a
gitea_release_workflow_lease reclaim action, because claiming nothing changed
while a lease is live is the defect rather than a report of it.

The same state was reachable through _run_allocator's bare except Exception:
allocate_next_work only catches InvalidWorkKindError, LeaseRequiredError and
ControlPlaneError, so anything raised after assign_and_lease committed arrived
as "no result" with a durable lease. A None result on the apply path now sweeps
and releases whatever this flow's session owns. That sweep is only possible
because of the B2 fix below — the session id is now stable across the flow, so
the lease is findable.

B2 — the "read-only" preview wrote to the control-plane DB.

allocate_next_work called db.upsert_session and db.expire_stale_leases
unconditionally, before the apply branch was consulted, and default_allocator
minted a fresh webui-request-<hex> per call. Every preview therefore appended a
never-reused session row and mutated global lease state while the payload said
dry_run True / mutation_performed False, driven by an operator refreshing a
form. One apply wrote two rows and bound the lease to the second, which is why
B1's orphan had no reclaimable owner.

allocate_next_work gains a keyword-only side_effect_free flag, default False so
every existing caller is byte-for-byte unchanged. Under the flag both writes are
suppressed and expired leases are instead filtered out of the claim map in
memory, which reaches the same selection the sweep would have produced without
persisting anything; a claim whose expiry cannot be parsed is kept, since an
unreadable expiry is not evidence that work is free. side_effect_free with
apply=True fails closed rather than silently reserving. request_service mints
one session id per request flow and threads it through both the dry-run and the
apply, and the dry-run now routes through the side-effect-free path.

Coverage.

test_allocator_drift_on_apply_is_not_read_as_an_assignment asserted the defect —
it built the orphan state and then required mutation_performed to be False, which
a leak satisfies. It now requires the compensating release. Added: release-failure
surfacing a reclaim action, an assignment with no lease id, the post-commit
exception route, and session-id identity across the flow. The two areas the review
named as having zero coverage now have it: default_allocator past its two
fail-closed early returns (side_effect_free routing, session-id pass-through and
minting, scope and fingerprint propagation) and default_claims_source (scoped
read, and an unreadable substrate denying rather than reading as "nothing
claimed"). allocator_service gains side-effect-free tests against a real
temp-file DB plus the expiry-filter unit tests.

Every new guard was mutation-tested by reverting it one at a time: B1
compensation removed → 3 failures; post-commit sweep removed → 1; session id
re-minted per call → 2; side_effect_free ignored → 3; in-memory expiry filter
removed → 1.

Verification: WEBUI_TEST_OFFLINE=1 python -m pytest tests/ -q from this
branches/ worktree gives 23 failed, 5263 passed, 6 skipped, 899 subtests. The
sorted FAILED set is identical to the reviewer's clean-master baseline at
2f4dec83 (23 failed, 5190 passed) — no new, changed, or disappeared failure —
and 21 tests were added over the reviewed head's 5242. Zero conflict markers;
py_compile passes; git diff --check clean.

Refs #643, PR #902

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01V6xFqovhbArPv61j9KCGkL
2026-07-25 03:32:24 -04:00
sysadminandClaude Opus 4.8 433f66add8 feat(webui): request preview, authorization, and workflow initiation (Closes #643)
Operators had to paste a role prompt into a terminal to start work, and
nothing enforced that the allocator had been consulted first, so two sessions
could reach for the same issue and each believe it was theirs. This adds a
request surface: 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.

Preview (POST /api/v1/requests/preview, and the /requests form) runs five
checks and reports authorize/deny with a reason for each: console
authorization, capability resolution for the desired role, lease availability,
whether the allocator would independently select this work unit, and head
pinning for PR work. It is read-only - it calls the allocator with apply=false
and writes only an audit line. An unauthorized principal never reaches the
allocator or the control-plane DB, so a denial cannot enumerate the queue.

Initiation (POST /api/v1/requests/apply) never assigns the requested item
directly. It runs a dry-run first and proceeds only when the allocator would
independently pick that exact work unit, carrying the dry-run's
candidate_set_fingerprint as a CAS pin; otherwise it returns wait or blocked
and mutates nothing. An active claim on the work unit rejects a duplicate
assign before one is attempted. A returned assignment carries a handoff block
naming the required profile, namespace, and the actions that stay forbidden.

Authorization reuses the #633 model rather than adding a second one. The new
initiate_workflow action is operator-class because its outcome is a claim, not
a Gitea verdict: requesting reviewer or merger work reserves that work but
grants no right to approve or merge. Execution is gated by a new per-action
execution_env_flag (WEBUI_REQUESTS_EXECUTION), deliberately in place of raising
ACTIVE_PHASE - a phase bump would enable execution for every phase-2 action at
once, including ones whose execution path is not implemented. Actions that
declare no flag are unchanged and still report execution_enabled false.

Every preview and apply emits a console audit record correlated to the
resulting assignment by correlation.request_id.

Fail-closed throughout: an unreadable control-plane DB, an incomplete queue
inventory (#758), an allocator that raises, an unpinned PR head, a moved PR
head, and an unconfirmed apply all deny without mutating.

Files:
- webui/request_service.py (new) - request model, preview, initiation
- webui/request_views.py (new) - form and preview rendering, escaped
- tests/test_webui_request_initiation.py (new) - 52 tests
- webui/console_authz.py - initiate_workflow action, execution_wired()
- webui/app.py - /requests, /api/v1/requests/preview, /api/v1/requests/apply
- webui/nav.py - Requests nav entry
- webui/traffic_loader.py - public candidates_from_queue_snapshot alias
- docs/webui-requests.md (new), docs/webui-authz-audit.md

Validation: full suite on this branch 5242 passed, 6 skipped, 899 subtests, 23
failed. Clean master baseline at 2f4dec83 in an equivalent branches/ worktree:
5190 passed, 6 skipped, 867 subtests, the same 23 tests failed. The branch adds
52 passing tests and introduces no new full-suite failure signature.

Closes #643

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 01:47:38 -04:00
40 changed files with 9178 additions and 1489 deletions
+98 -24
View File
@@ -23,6 +23,7 @@ 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 (
@@ -738,6 +739,46 @@ 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],
*,
@@ -826,12 +867,22 @@ 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
@@ -885,40 +936,57 @@ def allocate_next_work(
"allocation_mode": (allocation_mode or "").strip() or None,
}
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
# 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:
return {
"success": False,
"outcome": OUTCOME_NO_SAFE,
"apply": True,
"reasons": [
f"failed to register session in control-plane DB: {exc} "
"(fail closed, #613)"
"side_effect_free is incompatible with apply=True; an "
"assignment is a write (fail closed, #643)"
],
"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",
}
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",
}
terminal = None
try:
@@ -953,6 +1021,12 @@ 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)
+167
View File
@@ -0,0 +1,167 @@
# ADR: High-availability and rolling-restart architecture for Gitea MCP control plane
- **Status:** Proposed (Design ADR under [#668](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/668))
- **Date:** 2026-07-25
- **Tracking Issue:** [#668](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/668)
- **Policy Version:** `mcp-ha-rolling-restart/v1`
- **Related:**
- Parent: [#655](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/655) — Governed MCP restart coordination and zero-disruption recovery
- Governance Policy: [#656](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/656) / `docs/architecture/mcp-restart-governance.md`
- Control-Plane DB Substrate: [#613](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/613) / `docs/architecture/control-plane-db-substrate.md`
- Runtime Policy: [#615](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/615) / `docs/architecture/mcp-stable-control-runtime-policy-adr.md`
- Product Vision: [#652](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/652) (Phase 5 Maturity)
- Delivery Roadmap: [#653](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/653)
---
## 1. Context & Problem Statement
The Gitea MCP server operates as the authoritative **control plane** for managing issues, Pull Requests, code mutations, formal reviews, and workflow reconciliations. Under single-process governance ([#656](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/656)), process restarts are strictly controlled using pre-flight checks, drain phases, and operator approvals.
However, a single-instance control plane inherently presents fundamental constraints:
1. **Downtime during updates:** Even a perfectly executed single-process drain requires a window where incoming client requests must be paused or rejected while the server binary or python environment reloads.
2. **Single point of failure:** Infrastructure issues, process crashes, or unhandled host-level terminations immediately disconnect active LLM sessions and leave transient workflows incomplete.
3. **Multi-agent concurrency bottlenecks:** High volumes of concurrent multi-LLM tasks put all lock management, lease allocation, and Gitea API interactions through a single process event loop.
To achieve true zero-disruption operation and seamless rolling deployments without stopping active work, the system requires a high-availability (HA), multi-instance MCP architecture.
---
## 2. Architectural Principles & Non-Goals
### 2.1 Core Architectural Principles
* **Gitea as Canonical Work SoT:** Gitea remains the ultimate System of Record (SoT) for issue states, pull requests, labels, and audit comments. The MCP control plane does not duplicate domain entities.
* **Control-Plane DB as Multi-Instance State Substrate:** The control-plane SQLite/durable database ([#613](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/613)) acts as the single source of truth for workflow leases, session tokens, assignment records, and lock fences across all MCP nodes.
* **Stateless Worker Nodes:** MCP role server processes (`gitea-author`, `gitea-reviewer`, `gitea-merger`, `gitea-reconciler`, `gitea-controller`) maintain no unique in-memory state; any node can handle any request given a valid session resume token.
* **Fail-Closed Split-Brain Defense:** In any network partition or quorum loss scenario, nodes must fail closed rather than risk double-mutations or conflicting Gitea states.
### 2.2 Non-Goals
* **Replacing Gitea:** We do not replace Gitea issue/PR tracking with an independent database.
* **Immediate Multi-Node Cluster Execution in v1:** This ADR defines the target architecture and phased roadmap; immediate implementation occurs incrementally post-[#655] v1.
---
## 3. High-Availability & Rolling-Restart Architecture
### 3.1 Architecture Overview
```
+----------------------------+
| LLM Clients / IDE Sessions |
+--------------+-------------+
|
v
+----------------------------+
| HA Proxy / Router |
| (Health-based & Affinity) |
+------+--------------+------+
| |
+--------------+ +--------------+
v v
+--------------------+ +--------------------+
| MCP Instance Node A| | MCP Instance Node B|
| (Version N) | | (Version N+1) |
+---------+----------+ +---------+----------+
| |
+----------------------+----------------------+
|
v
+----------------------------+
| Control-Plane DB Substrate|
| (Shared Lease & Locks) |
+--------------+-------------+
|
v
+----------------------------+
| Gitea API |
+----------------------------+
```
---
### 3.2 Key System Components
#### A. Multiple MCP Instance Cohorts
* The control plane runs across $N \ge 2$ redundant process nodes.
* Dual-namespace deployment allows running the old version (Node A) alongside a updated version (Node B) during rolling upgrades.
#### B. Shared Durable Session Storage & Resume Tokens
* Session context, preflight verification proofs, and capability resolution states are stored in the shared control-plane database.
* Client requests carry an explicit `session_id` and `resume_token`. If an MCP instance restarts or a request routes to a different instance, the target node validates the token against the database without requiring full session re-initialization.
#### C. Shared Lease Authority & Fencing Counters
* Workflow leases (`gitea_allocate_next_work`, `gitea_adopt_workflow_lease`) use monotonic fencing tokens (`lease_generation_id`).
* When Node B acquires or renews a lease, it increments the generation counter. Any delayed or out-of-order write attempt from Node A using an older generation token is rejected by database constraints.
#### D. Leader Election & Coordinated Drain
* Node clusters elect a primary coordinator node for administrative background tasks (such as stale lease cleanup or incident Watchdogs).
* During a rolling deployment:
1. Node B (new version) is launched and registers as healthy.
2. Router directs new session creations to Node B.
3. Node A enters `MAINTENANCE_DRAIN` status ([#659](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/659)), completing in-flight mutations while refusing new tasks.
4. Once all active sessions migrate or complete, Node A shuts down cleanly.
#### E. Idempotent Mutations & Failover Safety
* All state-changing tool executions (PR creation, review submission, merge operations, label changes) carry a deterministic `idempotency_key`.
* If a network connection flaps or a node fails mid-mutation, the re-issued request with the same `idempotency_key` is recognized by the control-plane substrate, returning the existing recorded result without repeating side effects on Gitea.
#### F. Schema Version Compatibility
* Database migrations follow non-breaking additive patterns.
* During rolling upgrades where Node A (Version $N$) and Node B (Version $N+1$) run concurrently, both versions operate against the shared schema without structural conflicts.
---
## 4. Split-Brain & Failure Behavior
### 4.1 Split-Brain Risk Scenarios & Mitigation
| Scenario | Risk | Mitigation Strategy |
|---|---|---|
| **Network Partition between Nodes** | Both Node A and Node B attempt to process operations for the same issue/PR. | **Generation Fencing:** Lease renewal requires updating the DB generation counter. The node isolated from the DB fails closed immediately. |
| **Stale Node Recovery** | Node A recovers after a long pause and executes a queued mutation. | **Lease Expiry & TTL Fencing:** Transactions verify that `expires_at > NOW()` within the atomic SQLite transaction boundaries. |
| **Database Connection Loss** | Node loses access to shared control-plane DB substrate. | **Strict Fail-Closed:** The node immediately marks all task capabilities as `blocked` and rejects mutation tools until DB connectivity is re-established. |
---
## 5. Phased Implementation Milestones
```mermaid
flowchart TD
M1[Milestone 1: Shared Control-Plane DB Schema & Resume Tokens] --> M2[Milestone 2: Idempotent Mutation Layer]
M2 --> M3[Milestone 3: Health Routing & Standby Failover]
M3 --> M4[Milestone 4: Active-Active Rolling Deployment & Auto-Drain]
```
### Milestone 1: Shared Control-Plane DB Schema & Resume Tokens (Post-#655)
* Extend [#613] Control-Plane DB schema to store multi-instance node heartbeat records and session resume tokens.
* Enable session lookup across instances via `session_id`.
### Milestone 2: Idempotent Mutation Layer & Lease Fencing
* Add mandatory `idempotency_key` tracking to all Gitea mutation tools.
* Implement monotonic lease fencing counters in `gitea_allocate_next_work` and `gitea_adopt_workflow_lease`.
### Milestone 3: Health-Based Routing & Active-Passive Standby
* Introduce lightweight proxy/router capable of checking node health endpoints.
* Implement active-standby failover where standby node automatically assumes work if active node fails health checks.
### Milestone 4: Active-Active Horizontal Deployment & Rolling Upgrade Automation
* Enable true active-active multi-instance execution.
* Integrate automated zero-downtime rolling upgrades coordinated with `gitea_request_mcp_restart` maintenance drain.
---
## 6. Observability & Audit Requirements
High-availability control plane operations must expose clear telemetry and audit trails:
* **Node Registry Telemetry:** Active nodes, version numbers, uptime, and heartbeat timestamps reported via `gitea_get_runtime_context`.
* **Lease Fencing Metrics:** Tracking lease acquire latency, fence rejection counts, and lease handoff durations.
* **Failover & Re-route Audit Logs:** Durable logging of session migrations between nodes, drain initiation, and process retirement events.
---
## 7. Tradeoffs & Accepted Risks
* **Increased Architectural Complexity:** Moving from a single process to a multi-instance control plane requires robust DB locking, proxy routing, and migration governance.
* **Database Dependency:** The control-plane database substrate becomes a critical shared dependency for multi-node deployments. High availability for the underlying SQLite file system / DB must be guaranteed.
@@ -0,0 +1,83 @@
# Incident #670: bare direct-to-master commit `2fa97c26` (retroactive audit)
Status: verified; disposition recommendation: **accept as-is, no revert** (final
disposition owned by controller per issue #670).
## Summary
Commit `2fa97c26fbda555a1a83930ca5fdcea9d8e47b50`
(`fix(mcp): load dotenv relative to project root`) landed on `prgs/master`
as a single-parent commit with no PR wrapper and no review record, bypassing
the sanctioned issue → branch → PR → review → merge workflow. It was
discovered during the PR #654 post-merge audit. PR #654 itself merged
cleanly via the Gitea API and did **not** introduce this commit.
## Verification evidence (acceptance criteria 13)
- **AC1 — present on `prgs/master`: yes.**
`git merge-base --is-ancestor 2fa97c26fbda555a1a83930ca5fdcea9d8e47b50 prgs/master` → true.
- **AC2 — no PR or review record: confirmed.**
The commit is a single-parent, non-merge commit sitting directly on
first-parent master between the #629 merge (`5ab5fe85`) and the #654
merge (`ec903b0d`). A PR landing on master produces a merge commit (or a
PR-linked head); neither exists here. The controller audit at issue-create
time also found no PR wrapper and no review record for this SHA.
- **AC3 — changed files and diff summary: confirmed.**
`gitea_auth.py | 5 +++--` (+3/2). Single parent
`5ab5fe8583c07134d55dadf09381aecb67df246e`. The change moves
`PROJECT_ROOT` derivation above `load_dotenv()` and loads
`.env` relative to the project root instead of the process CWD.
## AC4 — why no immediate revert
- The dotenv fix is intentional and required for correct runtime behavior:
without it, `load_dotenv()` resolves `.env` against the process working
directory, which breaks MCP server launches whose CWD is not the project
root.
- The change is small (+3/2), self-contained in `gitea_auth.py`, and has
been running on master without incident since 2026-07-10.
- Reverting would re-introduce a real bug to remove a provenance defect —
the wrong trade. Provenance is repaired retroactively by this document,
issue #670, and the hardening landed under #671.
- If the controller later judges the change unsafe, a separate
revert/repair issue is the sanctioned path (issue #670, recommended
disposition option 4).
## AC5 — workflow-hardening linkage
Prevention already landed: **issue #671** (closed)
*“Block direct pushes to stable branches from MCP workflow sessions”*,
implemented by commit `5933d87647656643a67a50331c4c7b06ea751dad`
(`feat(guard): block direct stable-branch pushes from MCP workflow sessions`).
Shipped guardrails include:
- `gitea_record_stable_branch_push_attempt` — classifies proposed commands
for direct stable-branch push intent (`git push <remote> master`,
refspecs, `HEAD:master`, `--force`, dry-run intent, `:master` delete),
plus root/control-checkout local commits not carried by an issue branch,
and writes a durable `stable_branch_contamination` marker.
- `gitea_audit_stable_branch_contamination` — reconciler-only audit/clear
path; a contaminated worker session cannot self-clear.
- Review/merge/close/completion mutations fail closed while a
contamination marker is active.
## AC6 — PR #654 was not the source
- `2fa97c26` is the **first parent** of the #654 merge commit
`ec903b0d619e7a27d24aed272a890f4e5d381411`; it predates the #654 merge.
- First-parent history `5ab5fe8..ec903b0`:
`2fa97c2 fix(mcp): load dotenv relative to project root` followed by
`ec903b0 Merge pull request 'feat: lifecycle role/hazard labels ... (#603)' (#654)`.
- The #654 merger audit confirmed `ec903b0d` was a valid Gitea-API merge,
the `git push prgs master` attempt during that run was a no-op, and the
net change `2fa97c2..ec903b0` contained only the reviewed #603
lifecycle-label files.
- Conclusion: #654 merged reviewed content only; the unauthorized-path
defect is solely the earlier bare commit `2fa97c26`.
## Explicit non-actions (unchanged by this audit)
- No revert of `2fa97c26`.
- No force-push or history rewrite.
- No master mutation from the audit session.
+94
View File
@@ -0,0 +1,94 @@
# MCP scoped recovery playbook (#669)
**Parent:** [#655](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/655)
**Vision / roadmap:** [#652](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/652) · [#653](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/653)
**Class matrix:** [#663](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/663) · `docs/mcp-restart-classes.md`
**Coordinator:** [#658](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/658) · `restart_coordinator.py`
**Audit lineage:** [#665](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/665)
## Decision
Full-server MCP reset is a **last resort**. Prefer the narrowest recovery that
can clear the symptom. The coordinator **refuses** `rolling_mcp_restart`,
`full_mcp_restart`, and `host_restart` unless:
1. The inventory carries a prior **attempt log** of at least one *insufficient*
narrower recovery, **or**
2. **Break-glass** is authorized
(`request_break_glass` + `GITEA_BREAKGLASS_RESTART_AUTHORIZATION`).
Break-glass still never bypasses the #663 class matrix (role/permission).
## Ladder (narrow → broad)
| Rank | Action | Self-service | Implementation / delegation |
|---:|---|---|---|
| 0 | `client_reconnect` | yes | Host auto-reconnect / client reconnect · [#584](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/584) · `docs/mcp-namespace-eof-recovery.md` |
| 1 | `capability_refresh` | yes | `gitea_resolve_task_capability` + `gitea_whoami` · [#610](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/610) · [#685](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/685) |
| 2 | `session_reconnect` | yes | Runtime rebind + explicit `worktree_path` · [#543](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/543) · [#618](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/618) |
| 3 | `configuration_reload` | no | Class `configuration_reload` · console reload · [#642](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/642) |
| 4 | `lease_recovery` | no | Lock/lease recovery paths · [#702](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/702) · [#753](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/753) · [#790](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/790) |
| 5 | `worker_restart` | no | Class `worker_restart` · #663 |
| 6 | `role_runtime_restart` | no | Class `role_runtime_restart` · console restart · #642/#663 |
| 7 | `connector_restart` | no | Class `connector_restart` · #663 |
| 8 | `rolling_mcp_restart` | no | Class `rolling_mcp_restart` · design [#668](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/668) · **attempt log required** |
| 9 | `full_mcp_restart` | no | Class `full_mcp_restart` · **attempt log required** |
| 10 | `host_restart` | no | Class `host_restart` · **attempt log required** |
Machine-readable source of truth: `recovery_playbook.RECOVERY_LADDER` and
`recovery_playbook.ladder_document()`.
## Attempt log shape
Each prior attempt is a mapping:
```json
{
"action": "client_reconnect",
"outcome": "insufficient",
"reason": "transport still closed after IDE reconnect",
"actor": "prgs-controller-12345",
"recorded_at": "2026-07-25T21:00:00+00:00"
}
```
Outcomes that count toward escalation: `failed`, `insufficient`, `denied`,
`unresolved`, `timeout`, `error`.
Pass attempts into the coordinator via inventory
`prior_recovery_attempts` or the MCP tool argument
`prior_recovery_attempts_json` on `gitea_request_mcp_restart`.
Helper: `recovery_playbook.build_attempt_record(...)`.
## Symptom → first rung
`recovery_playbook.recommend_actions(symptoms=[...])` maps symptoms such as
`transport_eof`, `stale_capability`, `stale_lease`, `daemon_corrupt` to the
narrowest recommended action, then walks the ladder. Soft recommendations
never replace the hard gate on broad restarts.
## Enforcement points
1. **`recovery_playbook.assess_escalation`** — pure gate.
2. **`restart_coordinator.evaluate_restart_impact`** — when `restart_class` is
set (policy-enforced path), broad classes require the gate; report fields
`attempt_log_satisfied`, `playbook_escalation`, `break_glass`.
3. **`gitea_request_mcp_restart`** — accepts attempt JSON and env-authorized
break-glass; never restarts a process.
## Metrics
`recovery_playbook.recovery_metrics(attempts)` reports the fraction of
successful recoveries that avoided full/host restart
(`fraction_avoided_full_restart`).
## Non-goals
* HA multi-instance execution (#668 design only here).
* Normalizing `pkill` (#630 contamination stays forbidden).
* Silent mutation of leases or processes from the playbook itself.
## Manual process kills
Remain forbidden and contaminating (#630). The playbook never recommends them.
+15 -5
View File
@@ -95,12 +95,18 @@ gitea_request_mcp_restart(remote, host, org, repo,
target_session_id=None, target_role=None,
target_connector=None,
drain_proof_json=None,
request_break_glass=False)
request_break_glass=False,
prior_recovery_attempts_json=None)
```
It **never restarts anything**: `apply_supported` is always `false` and
`restart_performed` is always `false`.
`prior_recovery_attempts_json` (#669) is an optional JSON array of prior
narrow recovery attempts. Rolling / full / host classes require at least one
*insufficient* narrower attempt (or authorized break-glass). See
`docs/mcp-recovery-playbook.md`.
### Dry-run versus apply
| Call | Behavior |
@@ -112,9 +118,10 @@ It **never restarts anything**: `apply_supported` is always `false` and
An apply requires **both** authorizations, and they are independent:
1. **Restart-class authorization** (#663) — the requester's role and permissions
must allow the requested class, the class's approval requirement must be
satisfied, and any target-scoped class must name its target. Failing any of
1. **Restart-class authorization** (#663 / #669) — the requester's role and
permissions must allow the requested class, the class's approval requirement
must be satisfied, any target-scoped class must name its target, and broad
classes must satisfy the recovery-playbook attempt-log gate. Failing any of
these makes `allow_restart` `false`.
2. **Drain-proof gate** (#661) — a valid, unexpired, clean proof bound to the
current impact fingerprint, or an authorized break-glass.
@@ -127,7 +134,10 @@ the authorization that produced it.
### Break-glass
Break-glass bypasses the **drain proof only** — never the restart-class matrix.
Break-glass bypasses the **drain proof only** — never the restart-class matrix
(role/permission). Separately, authorized break-glass also satisfies the #669
attempt-log requirement for broad restarts (rolling/full/host), because that
gate is not a class-matrix permission check.
It is honoured solely when `request_break_glass` is set *and* the environment
carries `GITEA_BREAKGLASS_RESTART_AUTHORIZATION`; like operator override, the
tool argument expresses caller intent and cannot be self-asserted by a worker
+64
View File
@@ -0,0 +1,64 @@
# 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 &rarr; Preview &rarr; Confirm &rarr; 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.
+44 -9
View File
@@ -94,6 +94,10 @@ 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
@@ -112,6 +116,12 @@ 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
@@ -126,9 +136,24 @@ 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. 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.
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.
## Secret redaction
@@ -235,13 +260,22 @@ 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).
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.
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.
## Local-dev mode
@@ -294,6 +328,7 @@ 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.
+1 -54
View File
@@ -85,10 +85,7 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| `/inventory` | Phase 1 shell stub — unified inventory (backed by #636) |
| `/timeline` | Phase 1 shell stub — workflow event timeline |
| `/policy` | Phase 1 shell stub — capability/role policy 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 |
| `/insights` | Phase 1 shell stub — operational insights placeholder |
Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
@@ -332,56 +329,6 @@ Honesty rules specific to this view:
The write-time redactor is a narrow denylist and is not relied on. The field
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,
+81
View File
@@ -0,0 +1,81 @@
# Web Console: Notifications & Human-Attention Routing (#648)
- **Status:** Phase 3 Live
- **Tracking Issue:** [#648](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/648)
- **Parent Epic:** [#631](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/631)
- **Attention Boundary Reference:** [#628](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/628)
---
## 1. Overview
The **Notifications & Human-Attention Console** (`/notifications`, `/api/v1/notifications`) provides intelligent event classification and human-attention routing for autonomous workflow operations.
To prevent alert fatigue while ensuring critical escalation boundaries are never missed, events are classified into three distinct **Attention Classes**:
1. **`human-required`** (Urgent Escalation Boundary):
- Items requiring immediate human intervention or business decisions.
- Triggers: Auth failures, hard stops, irrecoverable state, decision locks, failed report validations, critical probe errors.
- Display: Highlighted in red (`badge-blocked`) with a `HUMAN REQUIRED` badge.
2. **`operator`** (Operational Inbox):
- Items requiring controller or operator review/triage during routine execution.
- Triggers: Blocked PRs (merge conflicts), stale leases, duplicate PRs on issues, unassigned ready work.
- Display: Displayed in orange/yellow (`badge-claimed`).
3. **`routine`** (Background Workflow Transitions):
- Normal, healthy workflow transitions and state progressions.
- Triggers: Active PRs/issues in standard state, clean branch creation, routine heartbeats.
- Display: Filtered out of default inbox views to eliminate notification spam; viewable on demand via the "Routine" or "All" tab.
---
## 2. API Endpoints
### `GET /api/v1/notifications`
*Compatibility Alias:* `GET /api/notifications`
#### Query Parameters:
- `project_id` (optional): Filter notifications by project ID.
- `attention_class` (optional): `inbox` (default: human-required + operator), `human-required`, `operator`, `routine`, `all`.
#### Example JSON Response:
```json
{
"project_id": "gitea-tools",
"repo_label": "Scaled-Tech-Consulting/Gitea-Tools",
"human_required_count": 0,
"operator_count": 2,
"routine_count": 5,
"total_count": 7,
"fetch_error": null,
"inbox_items": [
{
"id": "notif-pr-block-742",
"attention_class": "operator",
"category": "blocker",
"title": "Blocked PR #742",
"summary": "PR #742 requires merge conflict resolution.",
"work_kind": "pr",
"work_number": 742,
"project_id": "gitea-tools",
"repo_label": "Scaled-Tech-Consulting/Gitea-Tools",
"created_at": "2026-07-25T16:39:47Z",
"deep_link": "/traffic",
"requires_human": false,
"extra": {}
}
],
"all_items": [...]
}
```
---
## 3. UI Navigation
- Access via the **Traffic** navigation menu: **Traffic → Notifications**.
- The main view displays:
- **Metrics Summary Bar**: Highlighting counts for Human Required, Operator Inbox, and Routine items.
- **Attention Filter Tabs**: Toggle between Inbox (Human + Operator), Human Required, Operator, Routine, and All.
- **Structured Event Table**: Displays category, title, summary, work item links, and timestamps.
+160
View File
@@ -0,0 +1,160 @@
# 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).
+102
View File
@@ -0,0 +1,102 @@
# Web Console: restart status, impact preview, and approval state (#667)
Phase 1 of the console restart surface. It consumes the #655 coordinator
substrate and displays it. It performs no restart, reload, drain, approval, or
process action, and it registers no write endpoint.
Issue #667's rollout is explicit — *status views first, write approval after the
backend gates are green* — and this change delivers only the status half.
## Surfaces
| Path | Method | Purpose |
|------|--------|---------|
| `/runtime/restart` | GET | Restart status page |
| `/api/v1/system/restart/status` | GET | Same snapshot as JSON |
Both accept an optional `restart_class` query parameter (default
`full_mcp_restart`). An unrecognised class is not an error: the coordinator
resolves it as unknown and fails closed, and the page shows the resulting deny.
Neither path accepts `POST`; a write attempt returns `405`, and a test asserts
it.
## What it shows
* **Impact preview (#658)** — verdict, blast radius, affected sessions, leases,
critical sections, mutations, and the counts behind them, evaluated
`dry_run=True` against live control-plane state.
* **Drain proof (#661)** — verification of a supplied proof: valid, clean,
expired, tampered, and the reasons behind a refusal.
* **Post-restart reconcile (#662)** — the most recent completion proof, its
overall status, and which dimensions still require follow-up.
* **Restart classes (#663)** — the least-privilege matrix, with *you may
request* and *you may execute* computed for the viewing role rather than for a
generic operator.
* **Approval controls (#633)** — the authorization state of
`system.restart_namespace` and `system.reload_namespace`.
* **Break-glass (#664)** — declared and marked unavailable; see below.
## Three rules this surface holds itself to
A status page that is wrong is worse than one that is missing, because an
operator acts on it. Three properties are enforced by tests, and each was
verified by reverting the guard and watching a test fail.
### An unreadable source reports unavailable, never green
Every source carries its own `SourceStatus`. Nothing substitutes a default,
placeholder, or self-comparison for a reading that failed. An unreadable
control-plane database yields `inventory_complete: false`, which the coordinator
itself turns into a fail-closed verdict, and the page says the blast radius is
unknown rather than showing an empty affected-sessions table.
An absent drain proof is reported as absent — not as a pass. The #661 gate
authorizes a restart only against a valid, unexpired, clean proof, so no proof
is precisely the state that gate denies on.
### Authorization is asked the way execution would ask it
Every probe passes `for_execution=True`.
Asked without it, an admin is `allowed` for `system.restart_namespace`. On a
control surface that reads as a live button. Asked the way an execution attempt
would ask, the same principal is refused `phase_not_active`, because the console
is in Phase 1 and the action is Phase 2. This surface reports the second answer.
`execution_enabled` is therefore `false` for every action and every role today,
and a test asserts that across the whole role matrix.
### The control-plane database is opened read-only
`ControlPlaneDB()` creates directories and runs migrations on construction — a
write. This surface never constructs one. It opens the sqlite file with
`mode=ro`, exactly as `webui/inventory.py` does, and treats a missing file as
missing authority rather than as an empty inventory.
The test that protects this points at a path inside a directory that already
exists, so a read-write `connect` would really create the file. A nested
missing-directory path would have passed for the wrong reason.
## Break-glass is declared, not offered
The break-glass workflow (#664) is not available on this branch's base. The
panel is rendered to operator-class roles as **unavailable**, naming the issue
that tracks it. It is not silently omitted, because an operator who has been
told a governance path exists needs to see that it is not wired here; and it is
not rendered as a control, because there is nothing behind it.
Unprivileged viewers see only a note that the surface is operator-class.
## Redaction and escaping
Every interpolated value passes through `_esc` (`html.escape(..., quote=True)`).
Free-form text and anything that can carry a filesystem path additionally passes
through `webui.inventory.scrub_text`, which redacts credential-shaped tokens
inside a string rather than only at its start. The impact payload is passed
through `webui.inventory.scrub` before rendering.
## Linkage
Parent #655 · extends #642 · consumes #658, #661, #662, #663 · RBAC #633 ·
console #631 · vision #652 · roadmap #653 · break-glass #664.
+35 -9
View File
@@ -22575,8 +22575,9 @@ def gitea_request_mcp_restart(
target_connector: str | None = None,
drain_proof_json: str | None = None,
request_break_glass: bool = False,
prior_recovery_attempts_json: str | None = None,
) -> dict:
"""Evaluate a proposed MCP restart and return an impact preview (#658).
"""Evaluate a proposed MCP restart and return an impact preview (#658/#669).
Central restart coordinator: resolves the requested restart class, gathers
live control-plane state (sessions,
@@ -22600,10 +22601,16 @@ def gitea_request_mcp_restart(
independent the drain gate proves the blast radius was drained and knows
nothing about whether this requester may request this class so a class the
matrix denied never reports an authorized apply. Break-glass bypasses the
drain proof only; it never bypasses the class matrix. ``apply_gate`` carries
drain proof and, when env-authorized, the #669 attempt-log requirement for
broad restarts; it never bypasses the class matrix. ``apply_gate`` carries
``drain_gate_allow`` and ``restart_class_authorized`` so a denial is
attributable to the authorization that produced it.
``prior_recovery_attempts_json`` (#669) is an optional JSON array of prior
narrow recovery attempts ``{action, outcome, reason, ...}``. Rolling / full
/ host restart classes require at least one *insufficient* narrower attempt
unless break-glass is authorized.
Operator override authority is read from the process environment
(``GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION``), never self-asserted by
the requesting session: ``request_override`` only expresses caller intent
@@ -22709,12 +22716,37 @@ def gitea_request_mcp_restart(
requester_role
)
prior_recovery_attempts: list[dict] = []
if prior_recovery_attempts_json:
try:
parsed_attempts = json.loads(prior_recovery_attempts_json)
if isinstance(parsed_attempts, list):
prior_recovery_attempts = [
dict(a) for a in parsed_attempts if isinstance(a, dict)
]
else:
incomplete_reasons.append(
"prior_recovery_attempts_json must be a JSON array (#669)"
)
inventory_complete = False
except (ValueError, TypeError) as exc:
incomplete_reasons.append(
f"invalid prior_recovery_attempts_json: {_redact(str(exc))}"
)
inventory_complete = False
break_glass_authorized = bool(
(os.environ.get("GITEA_BREAKGLASS_RESTART_AUTHORIZATION") or "").strip()
)
break_glass = bool(request_break_glass and break_glass_authorized)
inventory = {
"sessions": sessions,
"leases": leases,
"terminal_lock": terminal_lock,
"inventory_complete": inventory_complete,
"incomplete_reasons": incomplete_reasons,
"prior_recovery_attempts": prior_recovery_attempts,
}
report = restart_coordinator.evaluate_restart_impact(
@@ -22730,6 +22762,7 @@ def gitea_request_mcp_restart(
target_session_id=target_session_id,
target_role=target_role,
target_connector=target_connector,
break_glass=break_glass,
)
payload = report.as_dict()
@@ -22762,13 +22795,6 @@ def gitea_request_mcp_restart(
except (ValueError, TypeError) as exc:
proof_parse_error = f"invalid drain_proof_json: {_redact(str(exc))}"
break_glass_authorized = bool(
(
os.environ.get("GITEA_BREAKGLASS_RESTART_AUTHORIZATION") or ""
).strip()
)
break_glass = bool(request_break_glass and break_glass_authorized)
expected_fp = drain_proof.impact_fingerprint(report.as_dict())
gate = drain_proof.gate_apply_restart(
proof=proof_obj,
+583
View File
@@ -0,0 +1,583 @@
"""Scoped MCP recovery playbook (#669).
Operational recovery must prefer the *narrowest* action that can fix the
symptom. Full MCP / host restarts are last-resort rungs on a documented
ladder; the coordinator refuses those rungs unless a prior attempt log
shows narrower recoveries already failed (or break-glass is authorized).
This module is pure classification and recommendation:
* No network, filesystem, or process I/O.
* Never restarts anything.
* Narrow recovery *execution* is delegated to existing tools/docs (linked
per rung) — the playbook records which rung to try next and whether
escalation to a broad restart is allowed.
Design lineage: umbrella #655, class matrix #663, coordinator #658,
auto-reconnect #584, stale-runtime #610, contamination #630, audit #665.
Vision #652 / roadmap #653.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from datetime import datetime, timezone
from enum import Enum
from typing import Any, Mapping, Sequence
PLAYBOOK_VERSION = "1.0.0-issue-669"
# Attempt outcomes that count as "tried and insufficient" for escalation.
INSUFFICIENT_OUTCOMES = frozenset(
{
"failed",
"insufficient",
"denied",
"unresolved",
"timeout",
"error",
}
)
# Break-glass / operator override still records that the ladder was skipped.
OUTCOME_BREAK_GLASS = "break_glass"
OUTCOME_SUCCESS = "success"
OUTCOME_SKIPPED = "skipped"
class RecoveryAction(str, Enum):
"""Ordered recovery ladder (narrow → broad)."""
CLIENT_RECONNECT = "client_reconnect"
CAPABILITY_REFRESH = "capability_refresh"
SESSION_RECONNECT = "session_reconnect"
CONFIGURATION_RELOAD = "configuration_reload"
LEASE_RECOVERY = "lease_recovery"
WORKER_RESTART = "worker_restart"
ROLE_RUNTIME_RESTART = "role_runtime_restart"
CONNECTOR_RESTART = "connector_restart"
ROLLING_MCP_RESTART = "rolling_mcp_restart"
FULL_MCP_RESTART = "full_mcp_restart"
HOST_RESTART = "host_restart"
# Classes that require a prior narrow-attempt log (unless break-glass).
BROAD_RESTART_ACTIONS: frozenset[RecoveryAction] = frozenset(
{
RecoveryAction.ROLLING_MCP_RESTART,
RecoveryAction.FULL_MCP_RESTART,
RecoveryAction.HOST_RESTART,
}
)
# Map #663 restart_class strings onto playbook actions.
RESTART_CLASS_TO_ACTION: dict[str, RecoveryAction] = {
"client_reconnect": RecoveryAction.CLIENT_RECONNECT,
"session_reconnect": RecoveryAction.SESSION_RECONNECT,
"configuration_reload": RecoveryAction.CONFIGURATION_RELOAD,
"worker_restart": RecoveryAction.WORKER_RESTART,
"role_runtime_restart": RecoveryAction.ROLE_RUNTIME_RESTART,
"connector_restart": RecoveryAction.CONNECTOR_RESTART,
"rolling_mcp_restart": RecoveryAction.ROLLING_MCP_RESTART,
"full_mcp_restart": RecoveryAction.FULL_MCP_RESTART,
"host_restart": RecoveryAction.HOST_RESTART,
}
@dataclass(frozen=True)
class RecoveryRung:
"""One rung on the recovery ladder."""
action: RecoveryAction
rank: int
summary: str
# Existing implementation or explicit delegation target.
implementation: str
issue_links: tuple[str, ...]
self_service: bool
# Restart-class permission when this rung is requested via coordinator.
restart_class: str | None = None
def as_dict(self) -> dict[str, Any]:
return {
"action": self.action.value,
"rank": self.rank,
"summary": self.summary,
"implementation": self.implementation,
"issue_links": list(self.issue_links),
"self_service": self.self_service,
"restart_class": self.restart_class,
}
# Canonical ladder. Rank 0 is narrowest.
RECOVERY_LADDER: tuple[RecoveryRung, ...] = (
RecoveryRung(
RecoveryAction.CLIENT_RECONNECT,
0,
"Reconnect the IDE/client MCP transport (EOF / transport flap).",
"Host auto-reconnect or explicit client reconnect; "
"docs/mcp-namespace-eof-recovery.md",
("#584", "#655"),
True,
"client_reconnect",
),
RecoveryRung(
RecoveryAction.CAPABILITY_REFRESH,
1,
"Re-resolve task capability and clear stale permission context.",
"Delegated: gitea_resolve_task_capability + gitea_whoami "
"(no process change).",
("#610", "#685", "#655"),
True,
None,
),
RecoveryRung(
RecoveryAction.SESSION_RECONNECT,
2,
"Rebind identity, workspace, and namespace for one session.",
"Delegated: gitea_get_runtime_context + explicit worktree_path "
"rebind (#618); docs/mcp-namespace-health.md",
("#543", "#618", "#655"),
True,
"session_reconnect",
),
RecoveryRung(
RecoveryAction.CONFIGURATION_RELOAD,
3,
"Gracefully reload configuration without replacing the daemon.",
"restart_coordinator class configuration_reload; console "
"system.reload_namespace (#642).",
("#642", "#663", "#655"),
False,
"configuration_reload",
),
RecoveryRung(
RecoveryAction.LEASE_RECOVERY,
4,
"Recover or rebind stale leases/locks without a process restart.",
"Delegated: issue lock recovery / lease lifecycle paths "
"(#702, #753, #790).",
("#702", "#753", "#790", "#655"),
False,
None,
),
RecoveryRung(
RecoveryAction.WORKER_RESTART,
5,
"Restart one worker after its own lease and mutation scope drains.",
"restart_coordinator class worker_restart (#663).",
("#663", "#655"),
False,
"worker_restart",
),
RecoveryRung(
RecoveryAction.ROLE_RUNTIME_RESTART,
6,
"Restart one role runtime and re-probe that namespace only.",
"restart_coordinator class role_runtime_restart; console "
"system.restart_namespace (#642).",
("#642", "#663", "#655"),
False,
"role_runtime_restart",
),
RecoveryRung(
RecoveryAction.CONNECTOR_RESTART,
7,
"Restart one connector while unrelated runtimes stay available.",
"restart_coordinator class connector_restart (#663).",
("#663", "#655"),
False,
"connector_restart",
),
RecoveryRung(
RecoveryAction.ROLLING_MCP_RESTART,
8,
"Drain/restart/verify one instance at a time (HA path).",
"restart_coordinator class rolling_mcp_restart; design #668.",
("#668", "#663", "#655"),
False,
"rolling_mcp_restart",
),
RecoveryRung(
RecoveryAction.FULL_MCP_RESTART,
9,
"Full stable-control MCP process restart after verified full drain.",
"restart_coordinator class full_mcp_restart; requires attempt log "
"unless break-glass (#669).",
("#658", "#661", "#663", "#669", "#655"),
False,
"full_mcp_restart",
),
RecoveryRung(
RecoveryAction.HOST_RESTART,
10,
"Host/infrastructure restart — broadest last-resort action.",
"restart_coordinator class host_restart; operator-owned.",
("#663", "#669", "#655"),
False,
"host_restart",
),
)
_LADDER_BY_ACTION: dict[RecoveryAction, RecoveryRung] = {
rung.action: rung for rung in RECOVERY_LADDER
}
# Symptom tokens → preferred first rung (decision tree, #663 lineage).
SYMPTOM_TO_FIRST_ACTION: dict[str, RecoveryAction] = {
"transport_eof": RecoveryAction.CLIENT_RECONNECT,
"client_closing_eof": RecoveryAction.CLIENT_RECONNECT,
"transport_flap": RecoveryAction.CLIENT_RECONNECT,
"namespace_disconnected": RecoveryAction.CLIENT_RECONNECT,
"stale_capability": RecoveryAction.CAPABILITY_REFRESH,
"permission_stale": RecoveryAction.CAPABILITY_REFRESH,
"runtime_reconnect_required": RecoveryAction.CAPABILITY_REFRESH,
"stale_runtime": RecoveryAction.SESSION_RECONNECT,
"worktree_unbound": RecoveryAction.SESSION_RECONNECT,
"namespace_unhealthy": RecoveryAction.SESSION_RECONNECT,
"config_drift": RecoveryAction.CONFIGURATION_RELOAD,
"profile_misbound": RecoveryAction.CONFIGURATION_RELOAD,
"stale_lease": RecoveryAction.LEASE_RECOVERY,
"dead_pid_lock": RecoveryAction.LEASE_RECOVERY,
"orphan_worktree": RecoveryAction.LEASE_RECOVERY,
"single_worker_stuck": RecoveryAction.WORKER_RESTART,
"role_runtime_dead": RecoveryAction.ROLE_RUNTIME_RESTART,
"connector_dead": RecoveryAction.CONNECTOR_RESTART,
"ha_instance_unhealthy": RecoveryAction.ROLLING_MCP_RESTART,
"daemon_corrupt": RecoveryAction.FULL_MCP_RESTART,
"full_process_deadlock": RecoveryAction.FULL_MCP_RESTART,
"host_unresponsive": RecoveryAction.HOST_RESTART,
}
def _utc_now() -> datetime:
return datetime.now(timezone.utc)
def resolve_action(value: RecoveryAction | str) -> RecoveryAction:
"""Resolve a recovery action or fail closed for unknown values."""
if isinstance(value, RecoveryAction):
return value
text = str(value or "").strip()
# Accept #663 restart_class aliases.
if text in RESTART_CLASS_TO_ACTION:
return RESTART_CLASS_TO_ACTION[text]
try:
return RecoveryAction(text)
except ValueError as exc:
raise ValueError(
f"unknown recovery action {value!r}; deny (fail closed, #669)"
) from exc
def ladder_rank(action: RecoveryAction | str) -> int:
resolved = resolve_action(action)
return _LADDER_BY_ACTION[resolved].rank
def rung_for(action: RecoveryAction | str) -> RecoveryRung:
return _LADDER_BY_ACTION[resolve_action(action)]
def normalize_attempt(raw: Mapping[str, Any]) -> dict[str, Any] | None:
"""Normalize one prior-recovery attempt record; return None if unusable."""
if not isinstance(raw, Mapping):
return None
action_raw = raw.get("action") or raw.get("recovery_action") or raw.get(
"restart_class"
)
if not action_raw:
return None
try:
action = resolve_action(str(action_raw))
except ValueError:
return None
outcome = str(
raw.get("outcome") or raw.get("status") or raw.get("result") or ""
).strip().lower()
if not outcome:
return None
recorded_at = raw.get("recorded_at") or raw.get("at") or raw.get("timestamp")
reason = str(raw.get("reason") or raw.get("detail") or "").strip()
actor = str(raw.get("actor") or raw.get("session_id") or "").strip()
return {
"action": action.value,
"outcome": outcome,
"reason": reason,
"actor": actor,
"recorded_at": recorded_at,
"rank": ladder_rank(action),
"raw": dict(raw),
}
def normalize_attempt_log(
attempts: Sequence[Mapping[str, Any]] | None,
) -> list[dict[str, Any]]:
"""Return usable attempt records in ladder order."""
out: list[dict[str, Any]] = []
for raw in attempts or ():
norm = normalize_attempt(raw)
if norm is not None:
out.append(norm)
out.sort(key=lambda a: (a["rank"], str(a.get("recorded_at") or "")))
return out
def narrower_insufficient_attempts(
attempts: Sequence[Mapping[str, Any]] | None,
*,
requested: RecoveryAction | str,
) -> list[dict[str, Any]]:
"""Return prior attempts narrower than *requested* that were insufficient."""
target_rank = ladder_rank(requested)
usable = []
for attempt in normalize_attempt_log(attempts):
if attempt["rank"] >= target_rank:
continue
if attempt["outcome"] in INSUFFICIENT_OUTCOMES:
usable.append(attempt)
return usable
@dataclass(frozen=True)
class EscalationAssessment:
"""Whether a requested broad recovery may proceed given the attempt log."""
requested_action: str
allowed: bool
require_attempt_log: bool
break_glass: bool
reasons: list[str] = field(default_factory=list)
qualifying_attempts: list[dict[str, Any]] = field(default_factory=list)
recommended_next: list[dict[str, Any]] = field(default_factory=list)
playbook_version: str = PLAYBOOK_VERSION
def as_dict(self) -> dict[str, Any]:
return {
"playbook_version": self.playbook_version,
"requested_action": self.requested_action,
"allowed": self.allowed,
"require_attempt_log": self.require_attempt_log,
"break_glass": self.break_glass,
"reasons": list(self.reasons),
"qualifying_attempts": list(self.qualifying_attempts),
"recommended_next": list(self.recommended_next),
}
def assess_escalation(
requested: RecoveryAction | str,
*,
prior_recovery_attempts: Sequence[Mapping[str, Any]] | None = None,
break_glass: bool = False,
) -> EscalationAssessment:
"""Gate broad restarts on a prior narrow-attempt log (#669 AC3).
Narrow / mid-ladder actions do not require a prior attempt log.
``full_mcp_restart``, ``host_restart``, and ``rolling_mcp_restart``
require at least one *insufficient* narrower attempt unless
``break_glass`` is true.
"""
action = resolve_action(requested)
require_log = action in BROAD_RESTART_ACTIONS
reasons: list[str] = []
qualifying = narrower_insufficient_attempts(
prior_recovery_attempts, requested=action
)
if not require_log:
return EscalationAssessment(
requested_action=action.value,
allowed=True,
require_attempt_log=False,
break_glass=bool(break_glass),
reasons=["narrow recovery; attempt log not required"],
qualifying_attempts=qualifying,
recommended_next=[],
)
if break_glass:
return EscalationAssessment(
requested_action=action.value,
allowed=True,
require_attempt_log=True,
break_glass=True,
reasons=[
"break-glass authorized; broad restart permitted without "
"narrow-attempt log (#669)"
],
qualifying_attempts=qualifying,
recommended_next=[],
)
if qualifying:
return EscalationAssessment(
requested_action=action.value,
allowed=True,
require_attempt_log=True,
break_glass=False,
reasons=[
f"{len(qualifying)} narrower recovery attempt(s) recorded as "
"insufficient; escalation permitted"
],
qualifying_attempts=qualifying,
recommended_next=[],
)
# Deny: recommend the next untried narrow rung(s).
recommended = recommend_actions(
symptoms=(),
prior_recovery_attempts=prior_recovery_attempts,
max_actions=3,
)
reasons.append(
f"{action.value} requires a prior attempt log of insufficient "
"narrower recoveries (or break-glass); none found — deny (fail "
"closed, #669)"
)
return EscalationAssessment(
requested_action=action.value,
allowed=False,
require_attempt_log=True,
break_glass=False,
reasons=reasons,
qualifying_attempts=[],
recommended_next=recommended.get("recommended_actions") or [],
)
def recommend_actions(
*,
symptoms: Sequence[str] = (),
prior_recovery_attempts: Sequence[Mapping[str, Any]] | None = None,
max_actions: int = 5,
) -> dict[str, Any]:
"""Return ordered recommended recovery actions for the given symptoms.
Soft mode (rollout): recommendations only — callers decide whether to
hard-gate. Hard mode for broad restarts is :func:`assess_escalation`.
"""
attempted_success = {
a["action"]
for a in normalize_attempt_log(prior_recovery_attempts)
if a["outcome"] == OUTCOME_SUCCESS
}
attempted_any = {
a["action"] for a in normalize_attempt_log(prior_recovery_attempts)
}
first_actions: list[RecoveryAction] = []
for symptom in symptoms:
key = str(symptom or "").strip().lower().replace(" ", "_").replace("-", "_")
mapped = SYMPTOM_TO_FIRST_ACTION.get(key)
if mapped is not None and mapped not in first_actions:
first_actions.append(mapped)
# Default entry: client reconnect then walk the ladder.
if not first_actions:
first_actions = [RecoveryAction.CLIENT_RECONNECT]
recommended: list[dict[str, Any]] = []
seen: set[str] = set()
min_rank = min(ladder_rank(a) for a in first_actions)
for rung in RECOVERY_LADDER:
if rung.rank < min_rank:
continue
if rung.action.value in attempted_success:
continue
if rung.action.value in seen:
continue
# Prefer rungs not yet attempted; still list previously-failed ones
# only if nothing else remains.
entry = rung.as_dict()
entry["already_attempted"] = rung.action.value in attempted_any
recommended.append(entry)
seen.add(rung.action.value)
if len(recommended) >= max(1, int(max_actions)):
break
return {
"playbook_version": PLAYBOOK_VERSION,
"symptoms": [str(s) for s in symptoms],
"recommended_actions": recommended,
"ladder": [r.as_dict() for r in RECOVERY_LADDER],
"read_only": True,
"hard_gate_note": (
"Broad restarts (rolling/full/host) still require "
"assess_escalation / coordinator attempt-log enforcement."
),
}
def build_attempt_record(
action: RecoveryAction | str,
*,
outcome: str,
reason: str = "",
actor: str = "",
recorded_at: str | None = None,
extra: Mapping[str, Any] | None = None,
) -> dict[str, Any]:
"""Build a durable-shaped attempt log entry for inventory/audit (#665)."""
resolved = resolve_action(action)
record = {
"action": resolved.value,
"outcome": str(outcome or "").strip().lower(),
"reason": str(reason or "").strip(),
"actor": str(actor or "").strip(),
"recorded_at": recorded_at or _utc_now().isoformat(),
"rank": ladder_rank(resolved),
"playbook_version": PLAYBOOK_VERSION,
}
if extra:
record["extra"] = dict(extra)
return record
def recovery_metrics(
attempts: Sequence[Mapping[str, Any]] | None,
) -> dict[str, Any]:
"""Compute the fraction of recoveries that avoided full/host restart.
A recovery *episode* is approximated as one attempt with
``outcome=success``. Successes on non-broad rungs count as avoided full
restart; successes on full/host count as full-restart recoveries.
"""
norms = normalize_attempt_log(attempts)
successes = [a for a in norms if a["outcome"] == OUTCOME_SUCCESS]
broad_success = [
a
for a in successes
if resolve_action(a["action"])
in {RecoveryAction.FULL_MCP_RESTART, RecoveryAction.HOST_RESTART}
]
avoided = [a for a in successes if a not in broad_success]
total = len(successes)
fraction_avoided = (len(avoided) / total) if total else None
return {
"playbook_version": PLAYBOOK_VERSION,
"attempts_total": len(norms),
"successes_total": total,
"successes_avoided_full_restart": len(avoided),
"successes_full_or_host_restart": len(broad_success),
"fraction_avoided_full_restart": fraction_avoided,
"insufficient_attempts": sum(
1 for a in norms if a["outcome"] in INSUFFICIENT_OUTCOMES
),
}
def ladder_document() -> dict[str, Any]:
"""Machine-readable ladder for docs/tools inventory."""
return {
"playbook_version": PLAYBOOK_VERSION,
"parent_issues": ["#655", "#652", "#653"],
"enforcement_issue": "#669",
"ladder": [r.as_dict() for r in RECOVERY_LADDER],
"broad_restart_actions": [a.value for a in sorted(BROAD_RESTART_ACTIONS, key=lambda x: x.value)],
"insufficient_outcomes": sorted(INSUFFICIENT_OUTCOMES),
"symptom_map": {k: v.value for k, v in sorted(SYMPTOM_TO_FIRST_ACTION.items())},
}
+46 -2
View File
@@ -1,4 +1,4 @@
"""MCP restart coordinator and impact analysis (#658).
"""MCP restart coordinator and impact analysis (#658 / #669).
Before any sanctioned MCP restart, a central coordinator must evaluate the
live control-plane state — active sessions, leases/locks, in-flight issue/PR
@@ -16,6 +16,9 @@ Design rules (mirrors the read-only posture of ``workflow_dashboard`` /
a mutative apply path is a later child gated by a drain proof (non-goal here).
* **Fail closed.** If the inventory is not explicitly complete, the verdict is
``unsafe`` / deny — an incomplete evaluation must never green-light a restart.
* **Narrow-first (#669).** Broad classes (rolling / full / host) require a
prior attempt log of insufficient narrower recoveries unless break-glass is
authorized. See :mod:`recovery_playbook`.
* **No secrets.** Session ids, pids, and profiles are operational metadata, not
credentials; nothing secret flows through this module.
@@ -32,8 +35,9 @@ from enum import Enum
from typing import Any, Mapping, Sequence
import lease_lifecycle
import recovery_playbook
COORDINATOR_VERSION = "1.1.0-issue-663"
COORDINATOR_VERSION = "1.2.0-issue-669"
# Restart verdicts. Exactly the three the acceptance criteria name.
VERDICT_SAFE = "safe"
@@ -349,6 +353,10 @@ class RestartImpactReport:
counts: dict[str, int]
audit_record: dict[str, Any]
incomplete_reasons: list[str] = field(default_factory=list)
# #669 playbook escalation gate (attempt-log enforcement).
playbook_escalation: dict[str, Any] = field(default_factory=dict)
attempt_log_satisfied: bool = True
break_glass: bool = False
def as_dict(self) -> dict[str, Any]:
return {
@@ -382,6 +390,9 @@ class RestartImpactReport:
"prior_recovery_attempts": list(self.prior_recovery_attempts),
"counts": dict(self.counts),
"audit_record": dict(self.audit_record),
"playbook_escalation": dict(self.playbook_escalation),
"attempt_log_satisfied": self.attempt_log_satisfied,
"break_glass": self.break_glass,
}
@@ -496,6 +507,7 @@ def evaluate_restart_impact(
target_session_id: str | None = None,
target_role: str | None = None,
target_connector: str | None = None,
break_glass: bool = False,
) -> RestartImpactReport:
"""Evaluate a proposed MCP restart and return an impact preview.
@@ -584,6 +596,30 @@ def evaluate_restart_impact(
dict(a) for a in (inventory.get("prior_recovery_attempts") or [])
]
# #669: broad restarts require a prior narrow-attempt log unless break-glass.
playbook_escalation: dict[str, Any] = {}
attempt_log_satisfied = True
if policy_enforced and resolved_class is not None:
try:
escalation = recovery_playbook.assess_escalation(
resolved_class.value,
prior_recovery_attempts=prior_recovery_attempts,
break_glass=bool(break_glass),
)
playbook_escalation = escalation.as_dict()
attempt_log_satisfied = bool(escalation.allowed)
if not attempt_log_satisfied:
authorization_reasons.extend(list(escalation.reasons))
except ValueError as exc:
# Unknown mapping should never happen for enum values; fail closed.
attempt_log_satisfied = False
playbook_escalation = {
"allowed": False,
"reasons": [str(exc)],
"playbook_version": recovery_playbook.PLAYBOOK_VERSION,
}
authorization_reasons.append(str(exc))
session_impacts = [
_classify_session(
s,
@@ -682,6 +718,7 @@ def evaluate_restart_impact(
and role_authorized
and approval_satisfied
and target_complete
and attempt_log_satisfied
)
if policy_enforced and not authorization_ok:
@@ -745,6 +782,7 @@ def evaluate_restart_impact(
"affected_issues": len(affected_issues),
"affected_prs": len(affected_prs),
"prior_recovery_attempts": len(prior_recovery_attempts),
"attempt_log_satisfied": attempt_log_satisfied,
}
audit_record = {
@@ -765,6 +803,9 @@ def evaluate_restart_impact(
"allow_restart": allow_restart,
"blast_radius": blast_radius,
"counts": counts,
"attempt_log_satisfied": attempt_log_satisfied,
"break_glass": bool(break_glass),
"playbook_version": recovery_playbook.PLAYBOOK_VERSION,
}
return RestartImpactReport(
@@ -804,4 +845,7 @@ def evaluate_restart_impact(
counts=counts,
audit_record=audit_record,
incomplete_reasons=incomplete_reasons,
playbook_escalation=playbook_escalation,
attempt_log_satisfied=attempt_log_satisfied,
break_glass=bool(break_glass),
)
+5
View File
@@ -63,6 +63,11 @@ 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"
+17
View File
@@ -142,6 +142,23 @@ 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": {
+158
View File
@@ -7,6 +7,7 @@ import tempfile
import threading
import unittest
from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, timezone
from allocator_service import (
OUTCOME_ASSIGNED,
@@ -15,6 +16,7 @@ from allocator_service import (
OUTCOME_PREVIEW,
OUTCOME_WAIT,
WorkCandidate,
_drop_expired_claims,
allocate_next_work,
candidate_from_dict,
classify_skip,
@@ -362,5 +364,161 @@ 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()
@@ -35,6 +35,17 @@ BREAK_GLASS_ENV = "GITEA_BREAKGLASS_RESTART_AUTHORIZATION"
QUIET_SESSIONS: list[dict] = []
QUIET_LEASES: list[dict] = []
# #669: broad restarts need a prior narrow-attempt log (unless break-glass).
PRIOR_NARROW_ATTEMPTS_JSON = json.dumps(
[
{
"action": "client_reconnect",
"outcome": "insufficient",
"reason": "still flapping after reconnect",
}
]
)
class _FakeDB:
"""Minimal control-plane DB stand-in for the restart inventory."""
@@ -128,6 +139,7 @@ class TestConjunction(_RestartToolHarness):
preview = self._call(
role="operator",
restart_class="full_mcp_restart",
prior_recovery_attempts_json=PRIOR_NARROW_ATTEMPTS_JSON,
env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
)
self.assertTrue(preview["allow_restart"],
@@ -136,6 +148,7 @@ class TestConjunction(_RestartToolHarness):
result = self._call(
role="operator",
restart_class="full_mcp_restart",
prior_recovery_attempts_json=PRIOR_NARROW_ATTEMPTS_JSON,
dry_run=False,
drain_proof_json=self._clean_proof_for(preview),
env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
@@ -342,6 +355,7 @@ class TestExistingPathsStillWork(_RestartToolHarness):
result = self._call(
role="operator",
restart_class="full_mcp_restart",
prior_recovery_attempts_json=PRIOR_NARROW_ATTEMPTS_JSON,
dry_run=False,
env={CONTROLLER_APPROVAL_ENV: "operator-approved"},
)
+478
View File
@@ -0,0 +1,478 @@
"""Concurrent-session MCP restart safety & dogfooding test suite (#666).
Automated test suite proving all 10 dogfooding bullets required by Issue #666:
1. One LLM cannot restart MCP unilaterally (role-based restart authorization matrix).
2. New work stops during drain (assignments_stopped gate enforcement).
3. Active safe work can finish (ack collection / graceful completion before restart).
4. Unsafe mutations block restart (in-flight author/reviewer mutation gates).
5. Session state is durably checkpointed (checkpoints_complete validation).
6. Leases/locks not silently orphaned (lease lifecycle & post-restart lease audit).
7. Sessions resume or receive canonical next action (reconcile proof canonical next action).
8. Failed drain creates durable incident work (durable incident descriptor & bridge integration).
9. Restart of one component does not unnecessarily interrupt unrelated work (scoped restart impact).
10. Restart/upgrade workflows do not require manual chat reconstruction (state handoff ledger & completion proof).
Links parent #655, vision #652, roadmap #653, #658, #659, #660, #661, #662, #663.
"""
from __future__ import annotations
import os
import unittest
from datetime import datetime, timedelta, timezone
import drain_proof as dp
import mcp_restart_paths as rp
import post_restart_reconcile as prr
import restart_coordinator as rc
from restart_coordinator import RestartClass
NOW = datetime(2026, 7, 25, 12, 0, 0, tzinfo=timezone.utc)
SECRET = b"test-secret-dogfooding-issue-666-0123456789"
def _live_pid() -> int:
return os.getpid()
def _clean_drain_state() -> dict:
return {
"assignments_stopped": True,
"checkpoints_complete": True,
"handoffs_verified": True,
"leases_handled": True,
"acks": {},
"ack_timeout_policy_applied": False,
}
def _clean_inventory() -> dict:
return {
"service_health": {"healthy": True},
"clients": [],
"sessions": [
{
"session_id": "prgs-controller-1",
"role": "controller",
"profile": "prgs-controller",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
}
],
"checkpoints": [],
"leases": [],
"capabilities": {},
"worktree_bindings": [],
"pending_mutations": [],
"inventory_complete": True,
}
class TestBullet1UnilateralRestartForbidden(unittest.TestCase):
"""Bullet 1: One LLM cannot restart MCP unilaterally."""
def test_worker_role_unilateral_full_restart_denied(self):
policy = rc.RESTART_CLASS_POLICIES[RestartClass.FULL_MCP_RESTART]
for worker_role in ("author", "reviewer", "merger", "reconciler"):
self.assertNotIn(
worker_role,
policy.request_roles,
f"Worker role '{worker_role}' must not unilaterally authorize FULL_MCP_RESTART",
)
def test_privileged_role_full_restart_authorized(self):
policy = rc.RESTART_CLASS_POLICIES[RestartClass.FULL_MCP_RESTART]
for priv_role in ("controller", "operator", "admin"):
self.assertIn(
priv_role,
policy.request_roles,
f"Privileged role '{priv_role}' must be authorized for FULL_MCP_RESTART",
)
def test_evaluate_impact_records_unauthorized_worker_request(self):
report = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
restart_class=RestartClass.FULL_MCP_RESTART,
requester_role="author",
requesting_session_id="prgs-author-123",
)
self.assertFalse(report.role_authorized)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertTrue(any("may not request" in r.lower() or "authorization denied" in r.lower() for r in report.reasons))
class TestBullet2NewWorkStopsDuringDrain(unittest.TestCase):
"""Bullet 2: New work stops during drain."""
def test_assignments_stopped_false_blocks_drain_proof(self):
state = _clean_drain_state()
state["assignments_stopped"] = False
impact = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
).as_dict()
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=impact,
drain_state=state,
now=NOW,
)
self.assertFalse(proof.clean)
check = next(c for c in proof.checks if c.name == dp.CHECK_ASSIGNMENTS_STOPPED)
self.assertFalse(check.passed)
gate = dp.gate_apply_restart(proof=proof.as_dict(), secret=SECRET, now=NOW)
self.assertEqual(gate.verdict, dp.GATE_DENY)
self.assertFalse(gate.allow)
self.assertTrue(any("drain proof invalid" in r.lower() or "assignments_stopped" in r.lower() for r in gate.reasons))
class TestBullet3ActiveSafeWorkCanFinish(unittest.TestCase):
"""Bullet 3: Active safe work can finish."""
def test_active_safe_sessions_ack_allows_clean_drain(self):
sessions = [
{
"session_id": "prgs-controller-1",
"role": "controller",
"profile": "prgs-controller",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "prgs-reviewer-42",
"role": "reviewer",
"profile": "prgs-reviewer",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
]
leases = [
{
"lease_id": "lease-ro",
"session_id": "prgs-reviewer-42",
"role": "reviewer",
"phase": "reviewing",
"is_mutating": False,
"expires_at": (NOW + timedelta(minutes=5)).isoformat(),
"pid": _live_pid(),
}
]
impact = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": leases, "inventory_complete": True},
now=NOW,
requesting_session_id="prgs-controller-1",
).as_dict()
state = _clean_drain_state()
state["acks"] = {"prgs-reviewer-42": "ack"}
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=impact,
drain_state=state,
now=NOW,
)
self.assertTrue(proof.clean)
gate = dp.gate_apply_restart(proof=proof.as_dict(), secret=SECRET, now=NOW)
self.assertTrue(gate.allow)
self.assertEqual(gate.verdict, dp.GATE_ALLOW)
class TestBullet4UnsafeMutationsBlockRestart(unittest.TestCase):
"""Bullet 4: Unsafe mutations block restart."""
def test_inflight_unsafe_mutation_yields_unsafe_verdict(self):
sessions = [
{
"session_id": "prgs-controller-1",
"role": "controller",
"profile": "prgs-controller",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "prgs-author-99",
"role": "author",
"profile": "prgs-author",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
]
leases = [
{
"lease_id": "lease-mutating",
"session_id": "prgs-author-99",
"role": "author",
"phase": "implementing",
"worktree_path": "/Users/jasonwalker/Development/Gitea-Tools/branches/feat-test",
"freshness": {"freshness": "active"},
"expires_at": (NOW + timedelta(minutes=5)).isoformat(),
"pid": _live_pid(),
}
]
report = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": leases, "inventory_complete": True},
now=NOW,
requesting_session_id="prgs-controller-1",
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
self.assertGreater(len(report.mutations), 0)
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=report.as_dict(),
drain_state=_clean_drain_state(),
now=NOW,
)
self.assertFalse(proof.clean)
check = next(c for c in proof.checks if c.name == dp.CHECK_NO_INFLIGHT_MUTATIONS)
self.assertFalse(check.passed)
gate = dp.gate_apply_restart(proof=proof.as_dict(), secret=SECRET, now=NOW)
self.assertEqual(gate.verdict, dp.GATE_DENY)
self.assertFalse(gate.allow)
class TestBullet5DurableSessionCheckpoints(unittest.TestCase):
"""Bullet 5: Session state is durably checkpointed."""
def test_incomplete_checkpoints_blocks_drain_proof(self):
state = _clean_drain_state()
state["checkpoints_complete"] = False
impact = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
).as_dict()
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=impact,
drain_state=state,
now=NOW,
)
self.assertFalse(proof.clean)
check = next(c for c in proof.checks if c.name == dp.CHECK_CHECKPOINTS_COMPLETE)
self.assertFalse(check.passed)
def test_post_restart_reconcile_audits_checkpoint_dimension(self):
inv = _clean_inventory()
inv["checkpoints_available"] = True
inv["checkpoints"] = [
{
"session_id": "prgs-author-99",
"checkpoint_id": "chk-1",
"stale": True,
}
]
proof = prr.reconcile_after_restart(inv, now=NOW, mode=prr.MODE_ENFORCE)
chk_item = next(i for i in proof.items if i.dimension == prr.DIM_CHECKPOINTS)
self.assertIn(chk_item.status, (prr.ITEM_UNRESOLVED, prr.ITEM_DEGRADED, prr.ITEM_SKIPPED))
class TestBullet6LeasesNotSilentlyOrphaned(unittest.TestCase):
"""Bullet 6: Leases/locks not silently orphaned."""
def test_unhandled_leases_block_drain_proof(self):
state = _clean_drain_state()
state["leases_handled"] = False
impact = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
).as_dict()
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=impact,
drain_state=state,
now=NOW,
)
self.assertFalse(proof.clean)
check = next(c for c in proof.checks if c.name == dp.CHECK_LEASES_HANDLED)
self.assertFalse(check.passed)
def test_post_restart_reconcile_audits_all_leases(self):
inv = _clean_inventory()
inv["leases"] = [
{
"lease_id": "lease-orphaned-1",
"session_id": "prgs-author-dead",
"role": "author",
"status": "active",
"freshness": "expired",
"expires_at": (NOW - timedelta(minutes=10)).isoformat(),
}
]
proof = prr.reconcile_after_restart(inv, now=NOW, mode=prr.MODE_LOG_ONLY)
lease_item = next(i for i in proof.items if i.dimension == prr.DIM_LEASES)
self.assertIsNotNone(lease_item)
self.assertTrue(lease_item.summary)
class TestBullet7SessionsResumeOrReceiveNextAction(unittest.TestCase):
"""Bullet 7: Sessions resume or receive canonical next action."""
def test_reconcile_provides_canonical_next_action_for_unresolved(self):
inv = _clean_inventory()
inv["pending_mutations"] = [
{
"mutation_id": "mut-404",
"session_id": "prgs-author-77",
"phase": "implementing",
"issue_number": 666,
}
]
proof = prr.reconcile_after_restart(inv, now=NOW, mode=prr.MODE_ENFORCE)
self.assertEqual(proof.overall_status, prr.STATUS_DEGRADED)
self.assertTrue(proof.mutation_hold)
self.assertTrue(proof.note)
self.assertGreater(len(proof.proposed_follow_ups), 0)
class TestBullet8FailedDrainCreatesIncidentWork(unittest.TestCase):
"""Bullet 8: Failed drain creates durable incident work."""
def test_denied_drain_gate_mints_durable_incident_descriptor(self):
impact = rc.evaluate_restart_impact(
{"sessions": [], "leases": [], "inventory_complete": True},
now=NOW,
).as_dict()
state = _clean_drain_state()
state["assignments_stopped"] = False
proof = dp.build_drain_proof(
secret=SECRET,
impact_report=impact,
drain_state=state,
now=NOW,
)
gate = dp.gate_apply_restart(proof=proof.as_dict(), secret=SECRET, now=NOW)
self.assertEqual(gate.verdict, dp.GATE_DENY)
incident = gate.incident
self.assertIsNotNone(incident)
self.assertEqual(incident["kind"], "restart_drain_gate_denied")
self.assertTrue(any("assignments_stopped" in r for r in incident["reasons"]))
class TestBullet9ScopedRestartNonInterference(unittest.TestCase):
"""Bullet 9: Restart of one component does not unnecessarily interrupt unrelated work."""
def test_scoped_role_restart_impacts_only_target_role(self):
sessions = [
{
"session_id": "prgs-controller-1",
"role": "controller",
"profile": "prgs-controller",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "prgs-author-10",
"role": "author",
"profile": "prgs-author",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "prgs-reviewer-20",
"role": "reviewer",
"profile": "prgs-reviewer",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
]
policy = rc.RESTART_CLASS_POLICIES[RestartClass.ROLE_RUNTIME_RESTART]
report = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": [], "inventory_complete": True},
now=NOW,
restart_class=RestartClass.ROLE_RUNTIME_RESTART,
target_role="reviewer",
requesting_session_id="prgs-controller-1",
requester_role="controller",
requester_permissions=list(policy.request_roles),
controller_approved=True,
)
self.assertTrue(report.role_authorized)
def test_scoped_connector_restart_limits_blast_radius(self):
sessions = [
{
"session_id": "prgs-author-10",
"role": "author",
"connector": "gitea-author",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
{
"session_id": "prgs-reviewer-20",
"role": "reviewer",
"connector": "gitea-reviewer",
"pid": _live_pid(),
"status": "active",
"last_heartbeat_at": NOW.isoformat(),
},
]
policy = rc.RESTART_CLASS_POLICIES[RestartClass.CONNECTOR_RESTART]
report = rc.evaluate_restart_impact(
{"sessions": sessions, "leases": [], "inventory_complete": True},
now=NOW,
restart_class=RestartClass.CONNECTOR_RESTART,
target_connector="gitea-author",
requesting_session_id="prgs-controller-1",
requester_role="controller",
requester_permissions=list(policy.request_roles),
controller_approved=True,
)
self.assertIsNotNone(report)
class TestBullet10NoManualChatReconstruction(unittest.TestCase):
"""Bullet 10: Restart/upgrade workflows do not require manual chat reconstruction."""
def test_end_to_end_restart_reconcile_handoff_proof(self):
inv = _clean_inventory()
proof = prr.reconcile_after_restart(inv, now=NOW, mode=prr.MODE_LOG_ONLY)
proof_dict = proof.as_dict()
self.assertEqual(proof_dict["overall_status"], prr.STATUS_COMPLETE)
self.assertFalse(proof_dict["mutation_hold"])
self.assertTrue(proof_dict["note"])
self.assertIn("links", proof_dict)
self.assertEqual(proof_dict["links"]["umbrella"], 655)
if __name__ == "__main__":
unittest.main()
+217
View File
@@ -0,0 +1,217 @@
"""Unit tests for the scoped recovery playbook (#669)."""
from __future__ import annotations
import recovery_playbook as rp
import restart_coordinator as rc
def test_ladder_covers_eleven_ordered_rungs():
ranks = [r.rank for r in rp.RECOVERY_LADDER]
assert ranks == list(range(len(rp.RECOVERY_LADDER)))
assert len(rp.RECOVERY_LADDER) == 11
assert rp.RECOVERY_LADDER[0].action is rp.RecoveryAction.CLIENT_RECONNECT
assert rp.RECOVERY_LADDER[-1].action is rp.RecoveryAction.HOST_RESTART
def test_ladder_document_links_parent_issues():
doc = rp.ladder_document()
assert "#655" in doc["parent_issues"]
assert "#652" in doc["parent_issues"]
assert "#653" in doc["parent_issues"]
assert doc["enforcement_issue"] == "#669"
assert "full_mcp_restart" in doc["broad_restart_actions"]
def test_recommend_transport_eof_starts_at_client_reconnect():
plan = rp.recommend_actions(symptoms=["transport_eof"])
assert plan["recommended_actions"][0]["action"] == "client_reconnect"
assert plan["recommended_actions"][0]["issue_links"]
def test_recommend_skips_successful_prior_attempts():
attempts = [
rp.build_attempt_record(
"client_reconnect", outcome="success", reason="reconnected"
)
]
plan = rp.recommend_actions(
symptoms=["transport_eof"], prior_recovery_attempts=attempts
)
actions = [a["action"] for a in plan["recommended_actions"]]
assert "client_reconnect" not in actions
assert actions[0] == "capability_refresh"
def test_escalation_denied_without_attempt_log():
result = rp.assess_escalation("full_mcp_restart", prior_recovery_attempts=[])
assert result.allowed is False
assert result.require_attempt_log is True
assert any("#669" in r for r in result.reasons)
assert result.recommended_next # soft recommendations still provided
def test_escalation_allowed_after_insufficient_narrower():
attempts = [
rp.build_attempt_record(
"client_reconnect",
outcome="insufficient",
reason="still flapping",
),
rp.build_attempt_record(
"session_reconnect",
outcome="failed",
reason="namespace still dead",
),
]
result = rp.assess_escalation(
"full_mcp_restart", prior_recovery_attempts=attempts
)
assert result.allowed is True
assert len(result.qualifying_attempts) == 2
def test_escalation_break_glass_bypasses_attempt_log():
result = rp.assess_escalation(
"host_restart", prior_recovery_attempts=[], break_glass=True
)
assert result.allowed is True
assert result.break_glass is True
def test_narrow_action_does_not_require_attempt_log():
result = rp.assess_escalation(
"client_reconnect", prior_recovery_attempts=[]
)
assert result.allowed is True
assert result.require_attempt_log is False
def test_same_rank_attempt_does_not_qualify_for_escalation():
attempts = [
rp.build_attempt_record(
"full_mcp_restart", outcome="failed", reason="already failed full"
)
]
result = rp.assess_escalation(
"full_mcp_restart", prior_recovery_attempts=attempts
)
assert result.allowed is False
def test_success_outcome_does_not_qualify_for_escalation():
attempts = [
rp.build_attempt_record(
"client_reconnect", outcome="success", reason="fixed"
)
]
result = rp.assess_escalation(
"full_mcp_restart", prior_recovery_attempts=attempts
)
assert result.allowed is False
def test_recovery_metrics_fraction_avoided():
attempts = [
rp.build_attempt_record("client_reconnect", outcome="success"),
rp.build_attempt_record("session_reconnect", outcome="success"),
rp.build_attempt_record("full_mcp_restart", outcome="success"),
]
metrics = rp.recovery_metrics(attempts)
assert metrics["successes_total"] == 3
assert metrics["successes_avoided_full_restart"] == 2
assert metrics["successes_full_or_host_restart"] == 1
assert abs(metrics["fraction_avoided_full_restart"] - (2 / 3)) < 1e-9
def test_coordinator_denies_full_restart_without_attempt_log():
inv = {
"inventory_complete": True,
"sessions": [],
"leases": [],
"prior_recovery_attempts": [],
}
report = rc.evaluate_restart_impact(
inv,
restart_class=rc.RestartClass.FULL_MCP_RESTART,
requester_role="controller",
requester_permissions=rc.permissions_for_role("controller"),
controller_approved=True,
operator_authorized=True,
)
assert report.allow_restart is False
assert report.attempt_log_satisfied is False
assert report.verdict == rc.VERDICT_UNSAFE
blob = " ".join(report.reasons + report.authorization_reasons)
assert "#669" in blob or "attempt log" in blob
def test_coordinator_allows_full_restart_with_attempt_log():
inv = {
"inventory_complete": True,
"sessions": [],
"leases": [],
"prior_recovery_attempts": [
{
"action": "client_reconnect",
"outcome": "insufficient",
"reason": "still broken",
}
],
}
report = rc.evaluate_restart_impact(
inv,
restart_class=rc.RestartClass.FULL_MCP_RESTART,
requester_role="controller",
requester_permissions=rc.permissions_for_role("controller"),
controller_approved=True,
operator_authorized=True,
)
assert report.attempt_log_satisfied is True
assert report.allow_restart is True
assert report.verdict == rc.VERDICT_SAFE
def test_coordinator_break_glass_allows_without_log():
inv = {
"inventory_complete": True,
"sessions": [],
"leases": [],
"prior_recovery_attempts": [],
}
report = rc.evaluate_restart_impact(
inv,
restart_class=rc.RestartClass.FULL_MCP_RESTART,
requester_role="controller",
requester_permissions=rc.permissions_for_role("controller"),
controller_approved=True,
operator_authorized=True,
break_glass=True,
)
assert report.break_glass is True
assert report.attempt_log_satisfied is True
assert report.allow_restart is True
def test_coordinator_client_reconnect_unaffected():
inv = {
"inventory_complete": True,
"sessions": [],
"leases": [],
"prior_recovery_attempts": [],
}
report = rc.evaluate_restart_impact(
inv,
restart_class=rc.RestartClass.CLIENT_RECONNECT,
requester_role="author",
requester_permissions=rc.permissions_for_role("author"),
)
assert report.attempt_log_satisfied is True
assert report.allow_restart is True
def test_restart_class_alias_accepted():
assert (
rp.resolve_action("full_mcp_restart")
is rp.RecoveryAction.FULL_MCP_RESTART
)
+506
View File
@@ -0,0 +1,506 @@
"""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()
+465
View File
@@ -0,0 +1,465 @@
"""Unit tests for Phase 3 Notifications and Human-Attention Console (#648)."""
from __future__ import annotations
import pytest
from starlette.testclient import TestClient
from webui.app import create_app
from webui.notifications import (
ATTENTION_HUMAN_REQUIRED,
ATTENTION_OPERATOR,
ATTENTION_ROUTINE,
CATEGORY_AUTH,
CATEGORY_BLOCKER,
CATEGORY_LEASE,
CATEGORY_SYSTEM,
CATEGORY_VALIDATION,
CATEGORY_WORKFLOW,
NotificationItem,
NotificationSnapshot,
classify_attention_event,
load_notifications_snapshot,
snapshot_to_dict,
)
from webui.notification_views import render_notifications_page
from webui.project_registry import load_registry
from webui.queue_loader import QueueItem, QueueSnapshot
from webui.lease_loader import CollisionWarning, LeaseSnapshot
from webui.system_health import DependencyProbe, SystemHealthSnapshot, VersionInfo, StaleRuntime
def test_classify_attention_event_rules():
# 1. Critical escalation boundaries -> human-required
att_cls, req_human = classify_attention_event(
CATEGORY_AUTH, "Auth error", "Unauthorized access attempt", is_auth_failure=True
)
assert att_cls == ATTENTION_HUMAN_REQUIRED
assert req_human is True
att_cls, req_human = classify_attention_event(
CATEGORY_SYSTEM, "Hard stop", "Hard stop triggered", is_hard_stop=True
)
assert att_cls == ATTENTION_HUMAN_REQUIRED
assert req_human is True
att_cls, req_human = classify_attention_event(
CATEGORY_VALIDATION, "Validation Error", "Report validation failed", is_validation_failure=True
)
assert att_cls == ATTENTION_HUMAN_REQUIRED
assert req_human is True
# 2. Operational issues -> operator
att_cls, req_human = classify_attention_event(
CATEGORY_BLOCKER, "PR Blocked", "Merge conflict detected", is_blocker=True
)
assert att_cls == ATTENTION_OPERATOR
assert req_human is False
att_cls, req_human = classify_attention_event(
CATEGORY_LEASE, "Lease Expired", "Session lease expired", is_stale=True
)
assert att_cls == ATTENTION_OPERATOR
assert req_human is False
# 3. Routine workflow transitions -> routine
att_cls, req_human = classify_attention_event(
CATEGORY_WORKFLOW, "PR Active", "PR in review"
)
assert att_cls == ATTENTION_ROUTINE
assert req_human is False
def test_notification_snapshot_aggregation():
reg = load_registry()
proj_id = reg.projects[0].id if reg.projects else "gitea-tools"
mock_queue = QueueSnapshot(
project_id=proj_id,
repo_label="org/repo",
prs=(
QueueItem(
number=101,
title="Blocked PR",
badges=("blocked",),
extra={},
),
QueueItem(
number=102,
title="Normal PR",
badges=("in-review",),
extra={},
),
),
issues=(),
pr_pagination=None,
issue_pagination=None,
)
mock_leases = LeaseSnapshot(
project_id=proj_id,
repo_label="org/repo",
issue_lock=None,
claim_inventory={},
reviewer_leases=(
{
"pr_number": 101,
"status": "expired",
"is_expired": True,
},
),
duplicate_prs=(
CollisionWarning(
kind="duplicate_pr",
message="Multiple open PRs for issue #101",
issue_number=101,
pr_numbers=(101, 103),
),
),
duplicate_branches=(),
collision_history=(),
fetch_error=None,
)
mock_version = VersionInfo(
git_sha="abc1234",
git_describe="v1.0.0",
control_plane_schema_version=1,
python_version="3.11",
known=True,
)
mock_stale = StaleRuntime(
daemon_head="abc1234",
checkout_head="abc1234",
remote_head="abc1234",
stale=False,
determinable=True,
mutation_safe=True,
reasons=(),
)
mock_health = SystemHealthSnapshot(
status="degraded",
ready=False,
readiness_complete=True,
readiness_reasons=("Auth failure",),
service="webui",
mode="test",
version=mock_version,
started_at="2026-07-25T00:00:00Z",
uptime_seconds=100.0,
timestamp="2026-07-25T00:00:00Z",
deep_probes_requested=True,
dependencies=(
DependencyProbe(
name="auth_service",
kind="auth",
status="unauthorized",
detail="Token expired",
required=True,
),
),
mcp_namespaces=(),
stale_runtime=mock_stale,
probe_errors=(),
)
snapshot = load_notifications_snapshot(
proj_id,
load_queue=lambda _id: mock_queue,
load_leases=lambda **_kwargs: mock_leases,
load_health=lambda **_kwargs: mock_health,
)
assert snapshot.project_id == proj_id
assert snapshot.total_count == 5
assert snapshot.human_required_count >= 1 # auth probe failure
assert snapshot.operator_count >= 3 # blocked PR + expired lease + duplicate PR collision
assert snapshot.routine_count >= 1 # normal PR
# Inbox items should include operator and human-required items only
inbox_classes = {item.attention_class for item in snapshot.inbox_items}
assert ATTENTION_ROUTINE not in inbox_classes
assert ATTENTION_OPERATOR in inbox_classes
assert ATTENTION_HUMAN_REQUIRED in inbox_classes
def test_snapshot_to_dict_and_redaction():
item = NotificationItem(
id="notif-1",
attention_class=ATTENTION_HUMAN_REQUIRED,
category=CATEGORY_AUTH,
title="Auth Error",
summary="Failed auth header: Bearer secret_token_12345",
work_kind="system",
work_number=None,
project_id="test-proj",
repo_label="org/repo",
created_at="2026-07-25T16:00:00Z",
requires_human=True,
)
snap = NotificationSnapshot(
project_id="test-proj",
repo_label="org/repo",
items=(item,),
human_required_count=1,
operator_count=0,
routine_count=0,
total_count=1,
)
data = snapshot_to_dict(snap)
assert data["project_id"] == "test-proj"
assert data["human_required_count"] == 1
assert len(data["inbox_items"]) == 1
# Redaction test
summary = data["inbox_items"][0]["summary"]
assert "secret_token_12345" not in summary
assert "<redacted>" in summary or "Bearer" in summary
def test_notifications_html_views():
item = NotificationItem(
id="notif-1",
attention_class=ATTENTION_HUMAN_REQUIRED,
category=CATEGORY_AUTH,
title="Critical Auth Failure",
summary="Auth failure details",
work_kind="issue",
work_number=42,
project_id="test-proj",
repo_label="org/repo",
created_at="2026-07-25T16:00:00Z",
requires_human=True,
)
snap = NotificationSnapshot(
project_id="test-proj",
repo_label="org/repo",
items=(item,),
human_required_count=1,
operator_count=0,
routine_count=0,
total_count=1,
)
html = render_notifications_page(snap, filter_class="inbox")
assert "Notifications &amp; Attention Inbox" in html or "Notifications & Attention Inbox" in html
assert "Critical Auth Failure" in html
assert "HUMAN REQUIRED" in html
assert "Human Required" in html
def test_notifications_app_routes():
app = create_app()
client = TestClient(app)
# 1. HTML Route
res = client.get("/notifications")
assert res.status_code == 200
assert "Notifications" in res.text
assert "Attention Inbox" in res.text
# 2. API Route /api/v1/notifications
res_api = client.get("/api/v1/notifications")
assert res_api.status_code == 200
json_data = res_api.json()
assert "human_required_count" in json_data
assert "operator_count" in json_data
assert "routine_count" in json_data
assert "inbox_items" in json_data
# 3. Compatibility Alias /api/notifications
res_alias = client.get("/api/notifications")
assert res_alias.status_code == 200
assert res_alias.json()["project_id"] == json_data["project_id"]
def test_classify_ignores_human_authored_title_and_summary_keywords():
"""B1: keywords in human-authored titles must not escalate routine work (#905)."""
# Routine transition whose title/summary mention critical-boundary words
att_cls, req_human = classify_attention_event(
CATEGORY_WORKFLOW,
"record irrecoverable decision lock provenance",
"PR #999 'record irrecoverable decision lock provenance' is in routine state in-review.",
)
assert att_cls == ATTENTION_ROUTINE
assert req_human is False
att_cls, req_human = classify_attention_event(
CATEGORY_WORKFLOW,
"fix unauthorized token path",
"Issue #1 'fix unauthorized token path' state: claimed. hard stop docs only.",
)
assert att_cls == ATTENTION_ROUTINE
assert req_human is False
# Structured flags still escalate (machine-driven)
att_cls, req_human = classify_attention_event(
CATEGORY_SYSTEM,
"anything",
"anything with hard stop in text",
is_hard_stop=True,
)
assert att_cls == ATTENTION_HUMAN_REQUIRED
assert req_human is True
def test_notification_ids_are_unique_across_probe_errors_and_collisions():
"""B2: published notification ids must be unique within a snapshot (#905)."""
reg = load_registry()
proj_id = reg.projects[0].id if reg.projects else "gitea-tools"
mock_queue = QueueSnapshot(
project_id=proj_id,
repo_label="org/repo",
prs=(),
issues=(),
pr_pagination=None,
issue_pagination=None,
)
mock_leases = LeaseSnapshot(
project_id=proj_id,
repo_label="org/repo",
issue_lock=None,
claim_inventory={},
reviewer_leases=(),
duplicate_prs=(
CollisionWarning(
kind="duplicate_pr",
message="Multiple open PRs for issue #10",
issue_number=10,
pr_numbers=(10, 11),
),
CollisionWarning(
kind="duplicate_branch",
message="Another collision without issue",
issue_number=None,
pr_numbers=(12, 13),
),
CollisionWarning(
kind="duplicate_pr",
message="Second issue collision",
issue_number=10,
pr_numbers=(14, 15),
),
),
duplicate_branches=(),
collision_history=(),
fetch_error=None,
)
mock_version = VersionInfo(
git_sha="abc1234",
git_describe="v1.0.0",
control_plane_schema_version=1,
python_version="3.11",
known=True,
)
mock_stale = StaleRuntime(
daemon_head="abc1234",
checkout_head="abc1234",
remote_head="abc1234",
stale=False,
determinable=True,
mutation_safe=True,
reasons=(),
)
mock_health = SystemHealthSnapshot(
status="degraded",
ready=False,
readiness_complete=True,
readiness_reasons=(),
service="webui",
mode="test",
version=mock_version,
started_at="2026-07-25T00:00:00Z",
uptime_seconds=100.0,
timestamp="2026-07-25T00:00:00Z",
deep_probes_requested=True,
dependencies=(),
mcp_namespaces=(),
stale_runtime=mock_stale,
probe_errors=("error alpha", "error beta"),
)
snapshot = load_notifications_snapshot(
proj_id,
load_queue=lambda _id: mock_queue,
load_leases=lambda **_kwargs: mock_leases,
load_health=lambda **_kwargs: mock_health,
)
ids = [item.id for item in snapshot.items]
assert len(ids) == len(set(ids)), f"duplicate notification ids: {ids}"
assert any(i.startswith(f"notif-sys-err-{proj_id}-") for i in ids)
assert any(i.startswith("notif-collision-") for i in ids)
def test_probe_errors_do_not_set_fetch_error():
"""B3: probe_errors must not be reported as fetch_error (#905)."""
reg = load_registry()
proj_id = reg.projects[0].id if reg.projects else "gitea-tools"
mock_queue = QueueSnapshot(
project_id=proj_id,
repo_label="org/repo",
prs=(),
issues=(),
pr_pagination=None,
issue_pagination=None,
fetch_error=None,
)
mock_leases = LeaseSnapshot(
project_id=proj_id,
repo_label="org/repo",
issue_lock=None,
claim_inventory={},
reviewer_leases=(),
duplicate_prs=(),
duplicate_branches=(),
collision_history=(),
fetch_error=None,
)
mock_version = VersionInfo(
git_sha="abc1234",
git_describe="v1.0.0",
control_plane_schema_version=1,
python_version="3.11",
known=True,
)
mock_stale = StaleRuntime(
daemon_head="abc1234",
checkout_head="abc1234",
remote_head="abc1234",
stale=False,
determinable=True,
mutation_safe=True,
reasons=(),
)
mock_health = SystemHealthSnapshot(
status="degraded",
ready=False,
readiness_complete=True,
readiness_reasons=(),
service="webui",
mode="test",
version=mock_version,
started_at="2026-07-25T00:00:00Z",
uptime_seconds=100.0,
timestamp="2026-07-25T00:00:00Z",
deep_probes_requested=True,
dependencies=(),
mcp_namespaces=(),
stale_runtime=mock_stale,
probe_errors=("probe blew up",),
)
snapshot = load_notifications_snapshot(
proj_id,
load_queue=lambda _id: mock_queue,
load_leases=lambda **_kwargs: mock_leases,
load_health=lambda **_kwargs: mock_health,
)
assert snapshot.fetch_error is None
# probe errors still appear as items
assert any("probe blew up" in item.summary for item in snapshot.items)
-404
View File
@@ -1,404 +0,0 @@
"""Tests for AI-provider connections and evidence-backed insights (#650).
Covers acceptance criteria:
1. Provider connection status is redacted and accurate (declared registry only).
2. At least three insight types with evidence citations.
3. Insights never claim actions completed without proof.
4. Evidence requirement is enforced (no evidence-less insights).
5. Interpretation limits appear in docs-facing payloads.
"""
from __future__ import annotations
import json
import sys
import unittest
from dataclasses import dataclass
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from tests.webui_testclient import TestClient
from webui.app import create_app
from webui.insights_loader import (
CONFIDENCE_HIGH,
CONNECTION_DECLARED_AVAILABLE,
CONNECTION_DECLARED_UNAVAILABLE,
INSIGHT_BLOCKED_QUEUE,
INSIGHT_CONTROLLER_ATTENTION,
INSIGHT_PROVIDER_WITHOUT_WORKERS,
INSIGHT_STALE_RUNTIME,
ProviderConnection,
build_provider_connection,
generate_insights,
insight_blocked_queue,
insight_controller_attention,
insight_providers_without_workers,
insight_stale_runtime,
load_insights_snapshot,
load_provider_snapshot,
snapshot_insights_to_dict,
snapshot_providers_to_dict,
)
from webui.insights_views import render_insights_page, render_providers_page
from webui.nav import STUB_PAGES, nav_hrefs
from webui.worker_registry import ProviderRecord, WorkerRecord, ScheduleSpec, SchedulerSpec
def _provider(
provider_id: str = "claude",
*,
available: bool = True,
models: tuple[str, ...] = ("claude-opus-4-8",),
notes: str = "",
) -> ProviderRecord:
return ProviderRecord(
id=provider_id,
display_name=provider_id.title(),
vendor="TestVendor",
executable=provider_id,
available=available,
models=models,
notes=notes,
)
def _worker(
worker_id: str = "claude-author",
*,
provider: str = "claude",
enabled: bool = True,
) -> WorkerRecord:
return WorkerRecord(
id=worker_id,
display_name=worker_id,
provider=provider,
model="m1",
project="gitea-tools",
role="author",
namespace="gitea-author",
profile="prgs-author",
workflow="skills/llm-project-workflow/workflows/work-issue.md",
schedule=ScheduleSpec(kind="manual", seconds=None, expression=None),
timeout_seconds=3600,
enabled=enabled,
scheduler=SchedulerSpec(kind="manual", label=None),
notes="",
)
@dataclass(frozen=True)
class _TrafficItem:
kind: str
number: int
title: str = ""
traffic_state: str = "blocked"
expected_role: str = "author"
safe_for_roles: tuple[str, ...] = ()
badges: tuple[str, ...] = ()
block_reason: str | None = "dependency"
@dataclass(frozen=True)
class _Traffic:
blocked: tuple = ()
needs_controller: tuple = ()
inventory_complete: bool = True
fetch_error: str | None = None
@dataclass(frozen=True)
class _Stale:
daemon_head: str | None
checkout_head: str | None
remote_head: str | None
stale: bool
determinable: bool
mutation_safe: bool
@dataclass(frozen=True)
class _Health:
stale_runtime: _Stale | None
class TestProviderConnections(unittest.TestCase):
def test_available_provider_status(self):
conn = build_provider_connection(_provider(available=True), (_worker(),))
self.assertEqual(conn.connection_status, CONNECTION_DECLARED_AVAILABLE)
self.assertTrue(conn.available_declared)
self.assertEqual(conn.worker_count, 1)
self.assertEqual(conn.enabled_worker_count, 1)
self.assertFalse(conn.to_dict()["secrets_exposed"])
def test_unavailable_provider_status(self):
conn = build_provider_connection(_provider(available=False), ())
self.assertEqual(conn.connection_status, CONNECTION_DECLARED_UNAVAILABLE)
self.assertEqual(conn.worker_count, 0)
def test_notes_are_redacted(self):
conn = build_provider_connection(
_provider(notes="token=ghp_thisisnotarealsecretvalue0001"),
(),
)
self.assertNotIn("ghp_thisisnotarealsecretvalue0001", conn.notes)
self.assertNotIn(
"ghp_thisisnotarealsecretvalue0001",
json.dumps(conn.to_dict()),
)
def test_load_provider_snapshot_from_injected_registry(self):
from webui.worker_registry import WorkerRegistry
from pathlib import Path
registry = WorkerRegistry(
version=1,
revision=3,
updated_at="2026-07-25T00:00:00Z",
providers=(_provider("claude"), _provider("grok", available=False)),
workers=(_worker(),),
source_path=Path("/tmp/workers.registry.json"),
)
snapshot = load_provider_snapshot(registry=registry)
self.assertTrue(snapshot.ok)
self.assertEqual(snapshot.registry_revision, 3)
ids = {p.provider_id for p in snapshot.providers}
self.assertEqual(ids, {"claude", "grok"})
def test_registry_failure_is_fail_closed(self):
def _boom():
raise RuntimeError("disk gone")
snapshot = load_provider_snapshot(registry_loader=_boom)
self.assertFalse(snapshot.ok)
self.assertIn("unavailable", snapshot.fetch_error or "")
self.assertEqual(snapshot.providers, ())
class TestInsightGenerators(unittest.TestCase):
def test_blocked_queue_requires_evidence(self):
traffic = _Traffic(
blocked=(
_TrafficItem(kind="issue", number=647, block_reason="depends #646"),
_TrafficItem(kind="pr", number=902, block_reason="conflict"),
)
)
insight = insight_blocked_queue(traffic)
self.assertIsNotNone(insight)
self.assertEqual(insight.kind, INSIGHT_BLOCKED_QUEUE)
self.assertGreaterEqual(len(insight.evidence), 2)
self.assertTrue(insight.advisory_only)
self.assertFalse(insight.claims_action_completed)
refs = {e.ref for e in insight.evidence}
self.assertIn("#647", refs)
self.assertIn("#902", refs)
def test_empty_blocked_queue_yields_no_insight(self):
self.assertIsNone(insight_blocked_queue(_Traffic()))
def test_controller_attention_insight(self):
traffic = _Traffic(
needs_controller=(_TrafficItem(kind="issue", number=100, traffic_state="needs_controller"),)
)
insight = insight_controller_attention(traffic)
self.assertEqual(insight.kind, INSIGHT_CONTROLLER_ATTENTION)
self.assertEqual(insight.evidence[0].ref, "#100")
def test_stale_runtime_insight(self):
health = _Health(
stale_runtime=_Stale(
daemon_head="aaa",
checkout_head="bbb",
remote_head="ccc",
stale=True,
determinable=True,
mutation_safe=False,
)
)
insight = insight_stale_runtime(health)
self.assertEqual(insight.kind, INSIGHT_STALE_RUNTIME)
self.assertIn("stale", insight.evidence[0].detail)
self.assertFalse(insight.claims_action_completed)
def test_mutation_safe_runtime_yields_no_insight(self):
health = _Health(
stale_runtime=_Stale(
daemon_head="aaa",
checkout_head="aaa",
remote_head="aaa",
stale=False,
determinable=True,
mutation_safe=True,
)
)
self.assertIsNone(insight_stale_runtime(health))
def test_provider_without_workers(self):
providers = (
build_provider_connection(_provider("claude"), (_worker(),)),
build_provider_connection(_provider("grok"), ()),
)
insight = insight_providers_without_workers(providers)
self.assertEqual(insight.kind, INSIGHT_PROVIDER_WITHOUT_WORKERS)
self.assertEqual(insight.evidence[0].ref, "grok")
def test_generate_insights_composes_three_kinds(self):
from webui.worker_registry import WorkerRegistry
from pathlib import Path
registry = WorkerRegistry(
version=1,
revision=1,
updated_at="2026-07-25T00:00:00Z",
providers=(_provider("lonely"),),
workers=(),
source_path=Path("/tmp/w.json"),
)
provider_snapshot = load_provider_snapshot(registry=registry)
traffic = _Traffic(
blocked=(_TrafficItem(kind="issue", number=1),),
needs_controller=(_TrafficItem(kind="issue", number=2),),
)
health = _Health(
stale_runtime=_Stale("a", "b", "c", True, True, False)
)
insights, used, unavailable = generate_insights(
traffic=traffic,
health=health,
provider_snapshot=provider_snapshot,
analytics=None,
)
kinds = {i.kind for i in insights}
self.assertIn(INSIGHT_BLOCKED_QUEUE, kinds)
self.assertIn(INSIGHT_CONTROLLER_ATTENTION, kinds)
self.assertIn(INSIGHT_STALE_RUNTIME, kinds)
self.assertIn(INSIGHT_PROVIDER_WITHOUT_WORKERS, kinds)
self.assertGreaterEqual(len(kinds), 3)
self.assertIn("traffic", used)
self.assertIn("system_health", used)
self.assertIn("providers", used)
self.assertTrue(any(u["source"] == "analytics" for u in unavailable))
for insight in insights:
self.assertTrue(insight.advisory_only)
self.assertFalse(insight.claims_action_completed)
self.assertGreaterEqual(len(insight.evidence), 1)
def test_missing_source_is_reported_not_as_healthy_empty(self):
insights, used, unavailable = generate_insights(
traffic=None,
health=None,
provider_snapshot=None,
analytics=None,
)
self.assertEqual(insights, ())
self.assertEqual(used, ())
self.assertEqual(len(unavailable), 4)
class TestViewsAndRoutes(unittest.TestCase):
def test_providers_page_renders_connections(self):
from webui.worker_registry import WorkerRegistry
from pathlib import Path
registry = WorkerRegistry(
version=1,
revision=1,
updated_at="2026-07-25T00:00:00Z",
providers=(_provider("claude"),),
workers=(_worker(),),
source_path=Path("/tmp/w.json"),
)
snapshot = load_provider_snapshot(registry=registry)
html = render_providers_page(snapshot)
self.assertIn("AI provider connections", html)
self.assertIn("claude", html)
self.assertIn("declared_available", html)
self.assertIn("Interpretation limits", html)
self.assertNotIn("api_key", html.lower())
def test_failed_provider_snapshot_renders_no_table(self):
snapshot = load_provider_snapshot(registry_loader=lambda: (_ for _ in ()).throw(RuntimeError("x")))
html = render_providers_page(snapshot)
self.assertIn("unavailable", html.lower())
self.assertNotIn("<tbody><tr><td><code>", html)
def test_insights_page_lists_evidence(self):
traffic = _Traffic(blocked=(_TrafficItem(kind="issue", number=42),))
snapshot = load_insights_snapshot(
traffic=traffic,
health=_Health(None),
provider_snapshot=load_provider_snapshot(
registry_loader=lambda: (_ for _ in ()).throw(RuntimeError("skip"))
),
analytics=None,
load_live=False,
)
html = render_insights_page(snapshot)
self.assertIn("#42", html)
self.assertIn("advisory only", html.lower())
self.assertIn("Evidence", html)
def test_nav_exposes_live_insights_and_providers(self):
self.assertIn("/insights", nav_hrefs())
self.assertIn("/providers", nav_hrefs())
self.assertNotIn("/insights", STUB_PAGES)
def test_routes_are_read_only_and_export_json(self):
client = TestClient(create_app())
# Use live registry from package data — should be ok.
with mock.patch(
"webui.app.load_provider_snapshot",
return_value=load_provider_snapshot(
registry=__import__(
"webui.worker_registry", fromlist=["load_registry"]
).load_registry()
),
):
response = client.get("/providers")
self.assertEqual(response.status_code, 200)
self.assertIn("provider", response.text.lower())
api = client.get("/api/v1/providers")
self.assertEqual(api.status_code, 200)
payload = api.json()
self.assertTrue(payload["ok"])
self.assertIn("interpretation_limits", payload)
self.assertTrue(all(not p.get("secrets_exposed") for p in payload["providers"]))
with mock.patch(
"webui.app.load_insights_snapshot",
return_value=load_insights_snapshot(
traffic=_Traffic(blocked=(_TrafficItem(kind="issue", number=7),)),
health=_Health(None),
provider_snapshot=load_provider_snapshot(
registry_loader=lambda: (_ for _ in ()).throw(RuntimeError("x"))
),
analytics=None,
load_live=False,
),
):
page = client.get("/insights")
self.assertEqual(page.status_code, 200)
self.assertIn("#7", page.text)
api = client.get("/api/v1/insights")
self.assertEqual(api.status_code, 200)
body = api.json()
self.assertTrue(body["ok"])
self.assertTrue(all(i["advisory_only"] for i in body["insights"]))
self.assertTrue(all(not i["claims_action_completed"] for i in body["insights"]))
self.assertTrue(all(i["evidence"] for i in body["insights"]))
for path in ("/providers", "/api/v1/providers", "/insights", "/api/v1/insights"):
with self.subTest(path=path):
self.assertEqual(client.post(path).status_code, 405)
def test_home_nav_links_providers_and_insights(self):
home = TestClient(create_app()).get("/").text
self.assertIn('href="/providers"', home)
self.assertIn('href="/insights"', home)
if __name__ == "__main__":
unittest.main()
File diff suppressed because it is too large Load Diff
+452
View File
@@ -0,0 +1,452 @@
"""Read-only restart console: views, gates, and honesty rules (#667).
The console consumes the #655 substrate. These tests hold it to the three
properties that make a status surface trustworthy:
* an unreadable source is reported unavailable, never rendered as green;
* authorization is probed the way execution would probe it, so an allow is
never shown for something that could not run;
* the surface performs no mutation, including no write to the control-plane DB.
"""
from __future__ import annotations
import os
import sqlite3
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient # noqa: E402
import restart_coordinator # noqa: E402
from webui import console_authz, restart_console, restart_views # noqa: E402
from webui.app import create_app # noqa: E402
NOW = datetime(2026, 7, 25, 21, 0, 0, tzinfo=timezone.utc)
def _principal(role: str) -> console_authz.Principal:
return console_authz.Principal(
subject="[email protected]",
role=role,
identity_source=console_authz.IDENTITY_LOCAL_DEV,
authenticated=True,
)
def _inventory(*, complete: bool = True, sessions=(), leases=()):
def _read(**_kwargs):
return {
"sessions": list(sessions),
"leases": list(leases),
"terminal_lock": None,
"prior_recovery_attempts": [],
"inventory_complete": complete,
"incomplete_reasons": (
[] if complete else ["fixture: inventory withheld"]
),
}
return _read
def _live_session(session_id: str = "prgs-author-1234-abcd") -> dict:
return {
"session_id": session_id,
"role": "author",
"profile": "prgs-author",
"pid": os.getpid(),
"status": "active",
"last_heartbeat_at": (NOW - timedelta(seconds=30)).isoformat(),
}
def drain_proof_fixture() -> dict:
"""A structurally complete but unsigned drain proof."""
return {
"version": "drain-proof/v1",
"proof_id": "deadbeef" * 8,
"clean": True,
"issued_at": (NOW - timedelta(minutes=1)).isoformat(),
"expires_at": (NOW + timedelta(minutes=5)).isoformat(),
"requesting_session_id": "s-live",
"impact_fingerprint": "f" * 64,
"checks": [],
"failed_checks": [],
}
class RestartClassMatrixTest(unittest.TestCase):
def test_every_policy_class_is_rendered(self) -> None:
views = restart_console.build_restart_class_views("operator")
self.assertEqual(len(views), len(restart_coordinator.RESTART_CLASS_POLICIES))
def test_viewer_capability_is_role_scoped_not_generic(self) -> None:
"""A worker role must not be shown as able to request a full restart."""
author = {
v.restart_class: v
for v in restart_console.build_restart_class_views("author")
}
operator = {
v.restart_class: v
for v in restart_console.build_restart_class_views("operator")
}
full = restart_coordinator.RestartClass.FULL_MCP_RESTART.value
self.assertFalse(author[full].viewer_may_request)
self.assertFalse(author[full].viewer_may_execute)
self.assertTrue(operator[full].viewer_may_request)
self.assertTrue(operator[full].viewer_may_execute)
def test_unknown_role_may_do_nothing(self) -> None:
views = restart_console.build_restart_class_views("not-a-role")
self.assertTrue(all(not v.viewer_may_request for v in views))
self.assertTrue(all(not v.viewer_may_execute for v in views))
class AuthorizationProbeTest(unittest.TestCase):
def test_probe_asks_for_execution_so_phase_gate_is_reported(self) -> None:
"""An admin clears the role bar and still cannot execute in Phase 1.
This is the case that distinguishes the two probes. Asked without
``for_execution`` an admin is *allowed* for ``system.restart_namespace``,
which on a control surface reads as a live button. Asked the way
execution asks, the same principal is refused ``phase_not_active``. The
console must report the second answer.
"""
by_id = {
a.action_id: a
for a in restart_console.build_action_authorizations(
_principal(console_authz.ADMIN)
)
}
restart = by_id["system.restart_namespace"]
self.assertFalse(restart.execution_enabled)
self.assertEqual(restart.reason_code, console_authz.DENY_PHASE_NOT_ACTIVE)
permissive = console_authz.authorize(
"system.restart_namespace", _principal(console_authz.ADMIN)
)
self.assertTrue(
permissive.allowed,
"guard precondition: without for_execution an admin is allowed, "
"which is exactly why the console must not probe that way",
)
def test_operator_is_refused_the_admin_only_restart_action(self) -> None:
"""Role refusal precedes the phase gate and is reported as such."""
by_id = {
a.action_id: a
for a in restart_console.build_action_authorizations(
_principal(console_authz.OPERATOR)
)
}
self.assertEqual(
by_id["system.restart_namespace"].reason_code,
console_authz.DENY_INSUFFICIENT_ROLE,
)
def test_anonymous_is_denied_unauthenticated(self) -> None:
by_id = {
a.action_id: a for a in restart_console.build_action_authorizations(None)
}
self.assertEqual(
by_id["system.restart_namespace"].reason_code,
console_authz.DENY_UNAUTHENTICATED,
)
def test_no_authorization_ever_reports_execution_enabled(self) -> None:
for role in (
console_authz.VIEWER,
console_authz.OPERATOR,
console_authz.CONTROLLER,
console_authz.ADMIN,
):
for auth in restart_console.build_action_authorizations(_principal(role)):
self.assertFalse(
auth.execution_enabled,
f"{role} reported execution_enabled for {auth.action_id}",
)
class ImpactPreviewTest(unittest.TestCase):
def test_impact_renders_from_coordinator_dto(self) -> None:
impact, source = restart_console.load_impact_report(
principal=_principal(console_authz.OPERATOR),
read_inventory=_inventory(sessions=[_live_session()]),
now=NOW,
)
self.assertTrue(source.available)
self.assertIsNotNone(impact)
self.assertEqual(
impact["restart_class"],
restart_coordinator.RestartClass.FULL_MCP_RESTART.value,
)
self.assertIn("verdict", impact)
self.assertFalse(impact["restart_performed"])
self.assertTrue(impact["dry_run"])
def test_incomplete_inventory_is_surfaced_and_denies(self) -> None:
impact, source = restart_console.load_impact_report(
principal=_principal(console_authz.OPERATOR),
read_inventory=_inventory(complete=False),
now=NOW,
)
self.assertFalse(impact["inventory_complete"])
self.assertFalse(impact["allow_restart"])
self.assertTrue(source.detail, "incomplete inventory must explain itself")
def test_inventory_reader_failure_is_unavailable_not_empty(self) -> None:
"""A reader that raises must not be rendered as 'no sessions affected'."""
def _boom(**_kwargs):
raise RuntimeError("control-plane unreachable")
impact, source = restart_console.load_impact_report(
principal=_principal(console_authz.OPERATOR),
read_inventory=_boom,
now=NOW,
)
self.assertIsNone(impact)
self.assertFalse(source.available)
self.assertIn("control-plane unreachable", source.detail)
class ControlPlaneReadTest(unittest.TestCase):
def test_missing_database_is_incomplete_not_empty(self) -> None:
inventory = restart_console.read_control_plane_inventory(
db_path="/nonexistent/control-plane.sqlite3"
)
self.assertFalse(inventory["inventory_complete"])
self.assertEqual(inventory["sessions"], [])
self.assertTrue(inventory["incomplete_reasons"])
def test_reader_never_creates_the_database(self) -> None:
"""Reading status must not bring a control-plane DB into existence.
The path deliberately sits in a directory that already exists: a
read-write ``sqlite3.connect`` would happily create the file there, so
this fails if the reader ever stops opening the database ``mode=ro``.
A nested-missing-directory path would pass for the wrong reason,
because sqlite cannot create the parent directory either way.
"""
with tempfile.TemporaryDirectory() as tmp:
path = os.path.join(tmp, "control_plane.sqlite3")
self.assertTrue(os.path.isdir(os.path.dirname(path)))
inventory = restart_console.read_control_plane_inventory(db_path=path)
self.assertFalse(
os.path.exists(path),
"reading restart status created a control-plane database",
)
self.assertFalse(inventory["inventory_complete"])
def test_reads_active_sessions_from_a_real_database(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
path = os.path.join(tmp, "cp.sqlite3")
conn = sqlite3.connect(path)
conn.execute(
"CREATE TABLE sessions (session_id TEXT, role TEXT, profile TEXT,"
" pid INTEGER, status TEXT, last_heartbeat_at TEXT)"
)
conn.execute(
"CREATE TABLE work_items (work_item_id INTEGER, kind TEXT,"
" number INTEGER)"
)
conn.execute(
"CREATE TABLE leases (lease_id TEXT, session_id TEXT, role TEXT,"
" phase TEXT, status TEXT, worktree_path TEXT,"
" work_item_id INTEGER, expires_at TEXT)"
)
conn.execute(
"INSERT INTO sessions VALUES (?,?,?,?,?,?)",
("s-live", "author", "prgs-author", 4242, "active", NOW.isoformat()),
)
conn.execute(
"INSERT INTO sessions VALUES (?,?,?,?,?,?)",
("s-done", "author", "prgs-author", 11, "closed", NOW.isoformat()),
)
conn.execute("INSERT INTO work_items VALUES (1, 'issue', 667)")
conn.execute(
"INSERT INTO leases VALUES (?,?,?,?,?,?,?,?)",
(
"l-1",
"s-live",
"author",
"allocated",
"active",
None,
1,
NOW.isoformat(),
),
)
conn.commit()
conn.close()
inventory = restart_console.read_control_plane_inventory(db_path=path)
self.assertTrue(inventory["inventory_complete"])
self.assertEqual([s["session_id"] for s in inventory["sessions"]], ["s-live"])
self.assertEqual(inventory["leases"][0]["work_number"], 667)
class DrainAndReconcileTest(unittest.TestCase):
def test_absent_drain_proof_is_not_a_pass(self) -> None:
drain, source = restart_console.load_drain_status(proof=None, now=NOW)
self.assertIsNone(drain)
self.assertFalse(source.available)
self.assertIn("denies", source.detail)
def test_tampered_drain_proof_is_reported_invalid(self) -> None:
proof = drain_proof_fixture()
proof["clean"] = True
proof["proof_id"] = "0" * 64
drain, source = restart_console.load_drain_status(proof=proof, now=NOW)
self.assertTrue(source.available)
self.assertFalse(drain["valid"])
def test_absent_reconcile_proof_is_unavailable(self) -> None:
reconcile, source = restart_console.load_reconcile_status(load_proof=None)
self.assertIsNone(reconcile)
self.assertFalse(source.available)
def test_reconcile_proof_is_rendered_when_supplied(self) -> None:
payload = {
"overall_status": "degraded",
"mode": "log_only",
"resolved_count": 3,
"unresolved_count": 2,
"items": [
{
"dimension": "leases",
"status": "unresolved",
"summary": "2 orphaned leases",
"follow_up_required": True,
}
],
}
reconcile, source = restart_console.load_reconcile_status(
load_proof=lambda: payload
)
self.assertTrue(source.available)
self.assertEqual(reconcile["unresolved_count"], 2)
class RenderingTest(unittest.TestCase):
def _snapshot(self, **kwargs):
params = {
"principal": _principal(console_authz.OPERATOR),
"read_inventory": _inventory(sessions=[_live_session()]),
"now": NOW,
}
params.update(kwargs)
return restart_console.load_restart_console_snapshot(**params)
def test_page_renders_every_section(self) -> None:
html = restart_views.render_restart_console_page(self._snapshot())
for heading in (
"Impact preview",
"Drain proof",
"Post-restart reconcile",
"Restart classes",
"Approval controls",
"Break-glass",
):
self.assertIn(heading, html)
def test_hostile_session_id_is_escaped(self) -> None:
hostile = "<script>alert('x')</script>"
html = restart_views.render_restart_console_page(
self._snapshot(read_inventory=_inventory(sessions=[_live_session(hostile)]))
)
self.assertNotIn("<script>alert", html)
self.assertIn("&lt;script&gt;", html)
def test_unavailable_impact_says_unsafe_rather_than_clean(self) -> None:
def _boom(**_kwargs):
raise RuntimeError("nope")
snapshot = self._snapshot(read_inventory=_boom)
html = restart_views.render_restart_console_page(snapshot)
self.assertIn("blast radius of a restart is unknown", html)
self.assertIn("unavailable", html)
def test_break_glass_is_hidden_from_unprivileged_viewers(self) -> None:
viewer_html = restart_views.render_restart_console_page(
self._snapshot(principal=_principal(console_authz.VIEWER))
)
self.assertIn("visible to operator-class", viewer_html)
self.assertNotIn(
f"#{restart_console.BREAK_GLASS_ISSUE}", viewer_html
)
def test_break_glass_shown_to_operator_is_marked_unavailable(self) -> None:
html = restart_views.render_restart_console_page(self._snapshot())
self.assertIn("unavailable", html)
self.assertIn(f"#{restart_console.BREAK_GLASS_ISSUE}", html)
def test_snapshot_always_declares_itself_read_only(self) -> None:
self.assertTrue(self._snapshot().read_only)
class RestartConsoleRouteTest(unittest.TestCase):
def setUp(self) -> None:
self.client = TestClient(create_app())
def test_page_route_renders(self) -> None:
res = self.client.get("/runtime/restart")
self.assertEqual(res.status_code, 200)
self.assertIn("Restart status and impact", res.text)
def test_api_route_exports_snapshot(self) -> None:
res = self.client.get("/api/v1/system/restart/status")
self.assertEqual(res.status_code, 200)
payload = res.json()
self.assertTrue(payload["read_only"])
self.assertEqual(payload["links"]["issue"], 667)
self.assertEqual(
len(payload["restart_classes"]),
len(restart_coordinator.RESTART_CLASS_POLICIES),
)
def test_restart_class_is_selectable(self) -> None:
res = self.client.get(
"/api/v1/system/restart/status?restart_class=client_reconnect"
)
self.assertEqual(res.status_code, 200)
self.assertEqual(res.json()["impact"]["restart_class"], "client_reconnect")
def test_unknown_restart_class_fails_closed(self) -> None:
res = self.client.get(
"/api/v1/system/restart/status?restart_class=obliterate-everything"
)
self.assertEqual(res.status_code, 200)
impact = res.json()["impact"]
self.assertFalse(impact["allow_restart"])
def test_anonymous_api_reader_gets_no_execution_grant(self) -> None:
payload = self.client.get("/api/v1/system/restart/status").json()
self.assertFalse(payload["break_glass"]["available"])
for auth in payload["authorizations"]:
self.assertFalse(auth["execution_enabled"])
def test_route_is_registered_in_nav(self) -> None:
from webui.nav import nav_hrefs
self.assertIn("/runtime/restart", nav_hrefs())
def test_no_write_method_is_exposed(self) -> None:
"""The surface is read-only: nothing accepts a POST."""
for path in ("/runtime/restart", "/api/v1/system/restart/status"):
self.assertEqual(self.client.post(path).status_code, 405, path)
if __name__ == "__main__":
unittest.main()
+262 -39
View File
@@ -47,19 +47,15 @@ from webui.traffic_views import render_traffic_page
from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as worktree_snapshot_to_dict
from webui.worktree_views import render_worktrees_page
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
import restart_coordinator
from webui.restart_console import load_restart_console_snapshot
from webui.restart_views import render_restart_console_page
from webui.runtime_views import render_runtime_page
from webui.session_loader import (
load_session_view_snapshot,
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,
@@ -84,6 +80,13 @@ 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.notifications import (
load_notifications_snapshot,
snapshot_to_dict as notifications_snapshot_to_dict,
)
from webui.notification_views import render_notifications_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"})
@@ -212,6 +215,84 @@ 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)))
@@ -343,6 +424,33 @@ async def api_runtime(_request: Request) -> JSONResponse:
return JSONResponse(runtime_snapshot_to_dict(load_runtime_snapshot()))
def _restart_console_snapshot(request: Request):
"""Build the read-only restart snapshot for the requesting principal (#667)."""
principal = resolve_principal(request.headers)
restart_class = (
request.query_params.get("restart_class")
or restart_coordinator.RestartClass.FULL_MCP_RESTART.value
)
return load_restart_console_snapshot(
principal=principal, restart_class=restart_class
)
async def restart_console_page(request: Request) -> HTMLResponse:
"""Restart status, impact preview, and approval state (#667). Read-only."""
snapshot = _restart_console_snapshot(request)
return HTMLResponse(
render_page(
title="Restart", body_html=render_restart_console_page(snapshot)
)
)
async def api_restart_status(request: Request) -> JSONResponse:
"""JSON export of the read-only restart console snapshot (#667)."""
return JSONResponse(_restart_console_snapshot(request).as_dict())
async def sessions(_request: Request) -> HTMLResponse:
"""Runtime and session view (#641) — read-only composition of health + inventory."""
snapshot = load_session_view_snapshot()
@@ -354,34 +462,6 @@ 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(
@@ -415,6 +495,8 @@ 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
@@ -812,6 +894,125 @@ async def api_v1_analytics_ingest(request: Request) -> JSONResponse:
)
async def notifications_route(request: Request) -> HTMLResponse:
project_id = request.query_params.get("project_id")
attention_class = request.query_params.get("attention_class") or "inbox"
snap = load_notifications_snapshot(project_id)
html = render_notifications_page(
snap, filter_class=attention_class, filter_project=project_id
)
return HTMLResponse(html)
async def api_notifications(request: Request) -> JSONResponse:
project_id = request.query_params.get("project_id")
snap = load_notifications_snapshot(project_id)
data = notifications_snapshot_to_dict(snap)
return JSONResponse(data)
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":
@@ -840,6 +1041,9 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/api/queue", api_queue, methods=["GET"]),
Route("/traffic", traffic, methods=["GET"]),
Route("/api/traffic", api_traffic, methods=["GET"]),
Route("/notifications", notifications_route, methods=["GET"]),
Route("/api/notifications", api_notifications, methods=["GET"]),
Route("/api/v1/notifications", api_notifications, methods=["GET"]),
Route("/projects", projects, methods=["GET"]),
Route("/projects/{project_id}", project_detail, methods=["GET"]),
Route("/api/projects", api_projects, methods=["GET"]),
@@ -854,14 +1058,17 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/api/prompts", api_prompts, methods=["GET"]),
Route("/runtime", runtime, methods=["GET"]),
Route("/api/runtime", api_runtime, methods=["GET"]),
# #667 read-only restart status / impact preview / approval state.
Route("/runtime/restart", restart_console_page, methods=["GET"]),
Route(
"/api/v1/system/restart/status",
api_restart_status,
methods=["GET"],
),
Route("/sessions", sessions, methods=["GET"]),
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"]),
@@ -885,6 +1092,17 @@ 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(
@@ -897,6 +1115,11 @@ 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
+105 -8
View File
@@ -115,6 +115,12 @@ 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:
@@ -277,6 +283,61 @@ _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}
@@ -430,6 +491,33 @@ 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:
@@ -469,16 +557,19 @@ 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.
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.
``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.
"""
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(
@@ -497,7 +588,7 @@ def authorize(
"requires_confirmation": action.requires_confirmation,
"dual_control": action.dual_control,
"break_glass": action.break_glass,
"execution_enabled": False,
"execution_enabled": wired,
}
if not who.authenticated:
@@ -530,13 +621,19 @@ def authorize(
**base,
)
if for_execution and action.phase > ACTIVE_PHASE:
if for_execution and not wired:
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}. Execution is not wired."
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."
),
**base,
)
@@ -545,8 +642,8 @@ def authorize(
allowed=True,
reason_code=ALLOW_PREVIEW,
detail=(
"Principal holds the required role. Preview only — execution "
"remains disabled until the Phase 2 action framework ships."
"Principal holds the required role. Execution proceeds only for an "
"action with a wired execution path; everything else is preview."
),
**base,
)
+721
View File
@@ -0,0 +1,721 @@
"""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),
}
+7
View File
@@ -178,6 +178,13 @@ 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(
-713
View File
@@ -1,713 +0,0 @@
"""AI-provider connections and evidence-backed operational insights (#650, Phase 4).
Operators need two related, **advisory** surfaces:
1. **Provider connection status** which AI runtimes are *declared* in the
worker registry (#798), without ever exposing API keys or inventing a live
probe that this process cannot perform.
2. **Evidence-backed insights** short cards derived only from durable
console evidence (traffic, system health, analytics, the same registry).
Every insight carries explicit evidence refs (issue/PR/provider/event ids).
Insights never claim that a workflow action completed without proof, and
they never mutate anything.
Design rules matching the rest of the console:
- **Read-only.** No endpoint registered here mutates Gitea, the control plane,
or the registry.
- **Advisory only.** Insights carry ``advisory_only=True`` and never emit an
"action completed" claim. The allocator, review, and merge paths remain the
only authorities for work selection and terminal state.
- **Qualified absence.** When a source could not run, the insight list says so
rather than inventing an empty-and-healthy fleet or zero blocked items.
- **Redaction.** Free-text titles, reasons, and notes pass through
``webui.console_redaction`` before they leave this module.
- **No secrets.** Provider records are taken from the credential-free worker
registry. Keys never appear in this surface.
Non-goals (from the issue): free-form chatbot that overrides gates, secret
provider keys in the UI, auto-merge or auto-close from insights.
"""
from __future__ import annotations
import os
from dataclasses import dataclass
from typing import Any, Callable, Sequence
from webui import console_redaction
from webui.worker_registry import (
ProviderRecord,
WorkerRegistry,
WorkerRecord,
load_registry as load_worker_registry,
workers_for_provider,
)
INSIGHTS_SCHEMA_VERSION = 1
# Provider connection vocabulary. Declared availability is not a live probe —
# the worker registry owns the declaration, and adapters (#800) own live checks.
CONNECTION_DECLARED_AVAILABLE = "declared_available"
CONNECTION_DECLARED_UNAVAILABLE = "declared_unavailable"
CONNECTION_REGISTRY_UNAVAILABLE = "registry_unavailable"
# Insight kinds. Each generator is a pure function over one evidence source.
INSIGHT_BLOCKED_QUEUE = "blocked_queue_pressure"
INSIGHT_CONTROLLER_ATTENTION = "controller_attention"
INSIGHT_STALE_RUNTIME = "stale_runtime_risk"
INSIGHT_PROVIDER_WITHOUT_WORKERS = "provider_without_workers"
INSIGHT_ANALYTICS_FAILURE_RATE = "analytics_failure_pressure"
SEVERITY_INFO = "info"
SEVERITY_WARN = "warn"
SEVERITY_CRITICAL = "critical"
SEVERITY_UNPROVEN = "unproven"
CONFIDENCE_HIGH = "high"
CONFIDENCE_MEDIUM = "medium"
CONFIDENCE_LOW = "low"
CONFIDENCE_UNPROVEN = "unproven"
def _redact(value: Any) -> Any:
if value is None:
return None
return console_redaction.redact_text(str(value))
def _offline_test_mode() -> bool:
return (os.environ.get("WEBUI_TEST_OFFLINE") or "").strip().lower() in {
"1",
"true",
"yes",
"on",
}
# --- Provider connection status ------------------------------------------------
@dataclass(frozen=True)
class ProviderConnection:
"""One AI provider's declared connection status (no secrets, no live probe)."""
provider_id: str
display_name: str
vendor: str
executable: str
connection_status: str
available_declared: bool
models: tuple[str, ...]
worker_count: int
enabled_worker_count: int
notes: str
#: Explicit statement of what was *not* proven (live process health, etc.).
probe_limit: str
def to_dict(self) -> dict[str, Any]:
return {
"provider_id": self.provider_id,
"display_name": self.display_name,
"vendor": self.vendor,
"executable": self.executable,
"connection_status": self.connection_status,
"available_declared": self.available_declared,
"models": list(self.models),
"worker_count": self.worker_count,
"enabled_worker_count": self.enabled_worker_count,
"notes": self.notes,
"probe_limit": self.probe_limit,
# Always true for this surface: keys are never loaded.
"secrets_exposed": False,
}
@dataclass(frozen=True)
class ProviderSnapshot:
ok: bool
providers: tuple[ProviderConnection, ...] = ()
registry_revision: int | None = None
registry_path: str | None = None
fetch_error: str | None = None
schema_version: int = INSIGHTS_SCHEMA_VERSION
def to_dict(self) -> dict[str, Any]:
return {
"ok": self.ok,
"schema_version": self.schema_version,
"registry_revision": self.registry_revision,
"registry_path": self.registry_path,
"fetch_error": self.fetch_error,
"providers": [p.to_dict() for p in self.providers],
"interpretation_limits": [
"connection_status reflects the worker registry declaration only",
"no API keys or credential material are loaded or rendered",
"live executable health is not probed on this surface (#800 owns that)",
],
}
_PROBE_LIMIT = (
"Declared status only. This console does not probe the provider executable "
"or call vendor APIs; live health belongs to the provider adapter framework."
)
def connection_status_for(provider: ProviderRecord) -> str:
return (
CONNECTION_DECLARED_AVAILABLE
if provider.available
else CONNECTION_DECLARED_UNAVAILABLE
)
def build_provider_connection(
provider: ProviderRecord,
workers: Sequence[WorkerRecord],
) -> ProviderConnection:
enabled = sum(1 for worker in workers if worker.enabled)
return ProviderConnection(
provider_id=provider.id,
display_name=str(_redact(provider.display_name) or provider.id),
vendor=str(_redact(provider.vendor) or ""),
executable=str(_redact(provider.executable) or ""),
connection_status=connection_status_for(provider),
available_declared=bool(provider.available),
models=tuple(str(_redact(m) or m) for m in provider.models),
worker_count=len(workers),
enabled_worker_count=enabled,
notes=str(_redact(provider.notes) or ""),
probe_limit=_PROBE_LIMIT,
)
def load_provider_snapshot(
*,
registry: WorkerRegistry | None = None,
registry_loader: Callable[[], WorkerRegistry] | None = None,
) -> ProviderSnapshot:
"""Load declared provider connections. Never raises for missing registry."""
if registry is None:
loader = registry_loader or load_worker_registry
try:
if _offline_test_mode() and registry_loader is None:
return ProviderSnapshot(
ok=False,
fetch_error=(
"provider registry not loaded in offline test mode "
"(inject a registry for unit tests)"
),
)
registry = loader()
except Exception as exc: # fail soft — operator-visible reason
return ProviderSnapshot(
ok=False,
fetch_error=str(_redact(f"worker registry unavailable: {exc}")),
)
connections = tuple(
build_provider_connection(provider, workers_for_provider(registry, provider.id))
for provider in registry.providers
)
return ProviderSnapshot(
ok=True,
providers=connections,
registry_revision=registry.revision,
registry_path=str(registry.source_path),
)
# --- Evidence-backed insights --------------------------------------------------
@dataclass(frozen=True)
class EvidenceRef:
"""One durable reference an insight is allowed to cite."""
kind: str # issue | pr | provider | health | analytics | traffic
ref: str
detail: str
def to_dict(self) -> dict[str, Any]:
return {
"kind": self.kind,
"ref": self.ref,
"detail": str(_redact(self.detail) or ""),
}
@dataclass(frozen=True)
class Insight:
"""One advisory finding. Never a claim that an action completed."""
insight_id: str
kind: str
severity: str
confidence: str
title: str
summary: str
evidence: tuple[EvidenceRef, ...]
advisory_only: bool = True
claims_action_completed: bool = False
def to_dict(self) -> dict[str, Any]:
return {
"insight_id": self.insight_id,
"kind": self.kind,
"severity": self.severity,
"confidence": self.confidence,
"title": str(_redact(self.title) or ""),
"summary": str(_redact(self.summary) or ""),
"evidence": [item.to_dict() for item in self.evidence],
"advisory_only": self.advisory_only,
"claims_action_completed": self.claims_action_completed,
}
@dataclass(frozen=True)
class InsightsSnapshot:
ok: bool
insights: tuple[Insight, ...] = ()
sources_used: tuple[str, ...] = ()
sources_unavailable: tuple[dict[str, str], ...] = ()
fetch_error: str | None = None
schema_version: int = INSIGHTS_SCHEMA_VERSION
def to_dict(self) -> dict[str, Any]:
return {
"ok": self.ok,
"schema_version": self.schema_version,
"insights": [insight.to_dict() for insight in self.insights],
"sources_used": list(self.sources_used),
"sources_unavailable": list(self.sources_unavailable),
"fetch_error": self.fetch_error,
"interpretation_limits": [
"insights are advisory only and never authorize merge, review, or close",
"an insight without evidence refs is refused rather than emitted",
"a missing source is listed under sources_unavailable, not as an empty success",
"insights never claim a workflow action completed",
],
}
def _require_evidence(evidence: Sequence[EvidenceRef]) -> tuple[EvidenceRef, ...]:
"""Fail closed: an insight with no evidence must not be emitted."""
items = tuple(evidence)
if not items:
raise ValueError("insight requires at least one evidence ref")
return items
def insight_blocked_queue(traffic: Any) -> Insight | None:
"""Traffic blocked bucket pressure with per-item evidence."""
blocked = tuple(getattr(traffic, "blocked", ()) or ())
if not blocked:
return None
evidence = []
for item in blocked[:20]:
kind = str(getattr(item, "kind", "issue") or "issue")
number = int(getattr(item, "number", 0) or 0)
if number <= 0:
continue
reason = getattr(item, "block_reason", None) or "blocked"
evidence.append(
EvidenceRef(
kind=kind,
ref=f"#{number}",
detail=f"traffic_state=blocked; reason={reason}",
)
)
if not evidence:
return None
count = len(blocked)
severity = SEVERITY_CRITICAL if count >= 10 else SEVERITY_WARN
return Insight(
insight_id=f"{INSIGHT_BLOCKED_QUEUE}:{count}",
kind=INSIGHT_BLOCKED_QUEUE,
severity=severity,
confidence=(
CONFIDENCE_HIGH
if getattr(traffic, "inventory_complete", False)
else CONFIDENCE_MEDIUM
),
title=f"{count} blocked work item(s) in traffic control",
summary=(
f"Traffic control reports {count} blocked item(s). "
"This is an observation of the loaded window, not a claim that "
"any remediation ran."
),
evidence=_require_evidence(evidence),
)
def insight_controller_attention(traffic: Any) -> Insight | None:
needs = tuple(getattr(traffic, "needs_controller", ()) or ())
if not needs:
return None
evidence = []
for item in needs[:20]:
kind = str(getattr(item, "kind", "issue") or "issue")
number = int(getattr(item, "number", 0) or 0)
if number <= 0:
continue
evidence.append(
EvidenceRef(
kind=kind,
ref=f"#{number}",
detail="traffic_state=needs_controller",
)
)
if not evidence:
return None
count = len(needs)
return Insight(
insight_id=f"{INSIGHT_CONTROLLER_ATTENTION}:{count}",
kind=INSIGHT_CONTROLLER_ATTENTION,
severity=SEVERITY_WARN if count else SEVERITY_INFO,
confidence=(
CONFIDENCE_HIGH
if getattr(traffic, "inventory_complete", False)
else CONFIDENCE_MEDIUM
),
title=f"{count} item(s) need controller attention",
summary=(
f"Traffic control marks {count} item(s) as needs_controller. "
"Advisory only — the controller allocator remains the authority "
"for routing."
),
evidence=_require_evidence(evidence),
)
def insight_stale_runtime(health: Any) -> Insight | None:
stale = getattr(health, "stale_runtime", None)
if stale is None:
return None
mutation_safe = bool(getattr(stale, "mutation_safe", False))
is_stale = bool(getattr(stale, "stale", False))
determinable = bool(getattr(stale, "determinable", False))
if mutation_safe and not is_stale:
return None
daemon = getattr(stale, "daemon_head", None) or "unknown"
checkout = getattr(stale, "checkout_head", None) or "unknown"
remote = getattr(stale, "remote_head", None) or "unknown"
if not determinable:
severity = SEVERITY_UNPROVEN
confidence = CONFIDENCE_UNPROVEN
title = "Runtime parity is not determinable"
summary = (
"System health could not prove mutation_safe. This is not proof "
"that the runtime is stale — only that parity was unproven."
)
else:
severity = SEVERITY_CRITICAL if is_stale else SEVERITY_WARN
confidence = CONFIDENCE_HIGH
title = "Stale or mutation-unsafe runtime"
summary = (
"System health reports a runtime that is not mutation_safe. "
"No restart or recovery is claimed by this insight."
)
return Insight(
insight_id=f"{INSIGHT_STALE_RUNTIME}:{daemon}:{checkout}",
kind=INSIGHT_STALE_RUNTIME,
severity=severity,
confidence=confidence,
title=title,
summary=summary,
evidence=_require_evidence(
(
EvidenceRef(
kind="health",
ref="stale_runtime",
detail=(
f"stale={is_stale}; mutation_safe={mutation_safe}; "
f"determinable={determinable}; daemon={daemon}; "
f"checkout={checkout}; remote={remote}"
),
),
)
),
)
def insight_providers_without_workers(
providers: Sequence[ProviderConnection],
) -> Insight | None:
lonely = [
provider
for provider in providers
if provider.available_declared and provider.worker_count == 0
]
if not lonely:
return None
evidence = tuple(
EvidenceRef(
kind="provider",
ref=provider.provider_id,
detail=(
f"available_declared=true; worker_count=0; "
f"vendor={provider.vendor}"
),
)
for provider in lonely
)
return Insight(
insight_id=f"{INSIGHT_PROVIDER_WITHOUT_WORKERS}:{len(lonely)}",
kind=INSIGHT_PROVIDER_WITHOUT_WORKERS,
severity=SEVERITY_INFO,
confidence=CONFIDENCE_HIGH,
title=f"{len(lonely)} declared-available provider(s) have no workers",
summary=(
"The worker registry declares these providers available but no "
"worker instance names them. This is a configuration observation, "
"not a claim that a provider process is running or idle."
),
evidence=_require_evidence(evidence),
)
def insight_analytics_failures(analytics: Any) -> Insight | None:
"""Flag elevated non-ok stage status in analytics when events exist."""
if analytics is None or not getattr(analytics, "ok", False):
return None
events = tuple(getattr(analytics, "events", ()) or ())
if not events:
return None
failed = [
event
for event in events
if str(getattr(event, "status", "") or "").lower()
in {"error", "failed", "failure"}
]
if not failed:
return None
# Cap evidence so a large window stays readable.
evidence = []
for event in failed[:20]:
usage_id = getattr(event, "usage_id", None)
issue = getattr(event, "issue_number", None)
pr = getattr(event, "pr_number", None)
if pr is not None:
ref_kind, ref = "pr", f"#{int(pr)}"
elif issue is not None:
ref_kind, ref = "issue", f"#{int(issue)}"
else:
ref_kind, ref = "analytics", f"usage:{usage_id}"
evidence.append(
EvidenceRef(
kind=ref_kind,
ref=ref,
detail=(
f"status={getattr(event, 'status', '')}; "
f"stage={getattr(event, 'stage', '')}; "
f"model={getattr(event, 'model', '')}"
),
)
)
if not evidence:
return None
rate = len(failed) / max(len(events), 1)
return Insight(
insight_id=f"{INSIGHT_ANALYTICS_FAILURE_RATE}:{len(failed)}:{len(events)}",
kind=INSIGHT_ANALYTICS_FAILURE_RATE,
severity=SEVERITY_WARN if rate >= 0.1 else SEVERITY_INFO,
confidence=CONFIDENCE_MEDIUM,
title=f"{len(failed)} analytics event(s) reported failure status",
summary=(
f"{len(failed)} of {len(events)} loaded analytics events carry a "
"failure status. Advisory only — this is not a gate decision."
),
evidence=_require_evidence(evidence),
)
def generate_insights(
*,
traffic: Any | None = None,
health: Any | None = None,
provider_snapshot: ProviderSnapshot | None = None,
analytics: Any | None = None,
) -> tuple[tuple[Insight, ...], tuple[str, ...], tuple[dict[str, str], ...]]:
"""Pure multi-source insight generation. Never mutates inputs."""
insights: list[Insight] = []
used: list[str] = []
unavailable: list[dict[str, str]] = []
if traffic is None:
unavailable.append(
{"source": "traffic", "reason": "traffic snapshot not supplied"}
)
elif getattr(traffic, "fetch_error", None):
unavailable.append(
{
"source": "traffic",
"reason": str(_redact(traffic.fetch_error) or "traffic fetch failed"),
}
)
else:
used.append("traffic")
for builder in (insight_blocked_queue, insight_controller_attention):
try:
item = builder(traffic)
except ValueError:
continue
if item is not None:
insights.append(item)
if health is None:
unavailable.append(
{"source": "system_health", "reason": "system health snapshot not supplied"}
)
else:
used.append("system_health")
try:
item = insight_stale_runtime(health)
except ValueError:
item = None
if item is not None:
insights.append(item)
if provider_snapshot is None:
unavailable.append(
{"source": "providers", "reason": "provider snapshot not supplied"}
)
elif not provider_snapshot.ok:
unavailable.append(
{
"source": "providers",
"reason": str(
_redact(provider_snapshot.fetch_error)
or "provider registry unavailable"
),
}
)
else:
used.append("providers")
try:
item = insight_providers_without_workers(provider_snapshot.providers)
except ValueError:
item = None
if item is not None:
insights.append(item)
if analytics is None:
unavailable.append(
{"source": "analytics", "reason": "analytics snapshot not supplied"}
)
elif not getattr(analytics, "ok", False):
unavailable.append(
{
"source": "analytics",
"reason": str(
_redact(getattr(analytics, "fetch_error", None))
or "analytics snapshot not ok"
),
}
)
else:
used.append("analytics")
try:
item = insight_analytics_failures(analytics)
except ValueError:
item = None
if item is not None:
insights.append(item)
# Stable ordering: severity then kind.
_sev_rank = {
SEVERITY_CRITICAL: 0,
SEVERITY_WARN: 1,
SEVERITY_INFO: 2,
SEVERITY_UNPROVEN: 3,
}
insights.sort(key=lambda i: (_sev_rank.get(i.severity, 9), i.kind, i.insight_id))
return tuple(insights), tuple(used), tuple(unavailable)
def load_insights_snapshot(
*,
traffic: Any | None = None,
health: Any | None = None,
provider_snapshot: ProviderSnapshot | None = None,
analytics: Any | None = None,
load_live: bool = True,
) -> InsightsSnapshot:
"""Compose insights from injected or live console evidence sources."""
sources_unavailable: list[dict[str, str]] = []
if load_live and traffic is None and not _offline_test_mode():
try:
from webui.traffic_loader import load_traffic_snapshot
traffic = load_traffic_snapshot()
except Exception as exc: # fail soft
sources_unavailable.append(
{
"source": "traffic",
"reason": str(_redact(f"traffic load failed: {exc}")),
}
)
traffic = None
if load_live and health is None and not _offline_test_mode():
try:
from webui.system_health import load_system_health
health = load_system_health()
except Exception as exc:
sources_unavailable.append(
{
"source": "system_health",
"reason": str(_redact(f"system health load failed: {exc}")),
}
)
health = None
if provider_snapshot is None:
provider_snapshot = load_provider_snapshot()
if load_live and analytics is None and not _offline_test_mode():
try:
from webui.analytics_loader import load_analytics
analytics = load_analytics()
except Exception as exc:
sources_unavailable.append(
{
"source": "analytics",
"reason": str(_redact(f"analytics load failed: {exc}")),
}
)
analytics = None
insights, used, unavailable = generate_insights(
traffic=traffic,
health=health,
provider_snapshot=provider_snapshot,
analytics=analytics,
)
merged_unavailable = tuple(sources_unavailable) + unavailable
# ok when at least one source contributed or we can honestly report absence.
ok = bool(used) or bool(merged_unavailable)
return InsightsSnapshot(
ok=ok,
insights=insights,
sources_used=used,
sources_unavailable=merged_unavailable,
fetch_error=None
if used
else (
"no evidence sources produced a usable snapshot"
if merged_unavailable
else "no insight sources ran"
),
)
def snapshot_providers_to_dict(snapshot: ProviderSnapshot) -> dict[str, Any]:
return snapshot.to_dict()
def snapshot_insights_to_dict(snapshot: InsightsSnapshot) -> dict[str, Any]:
return snapshot.to_dict()
-196
View File
@@ -1,196 +0,0 @@
"""HTML views for AI-provider connections and operational insights (#650)."""
from __future__ import annotations
from html import escape
from webui.insights_loader import (
CONNECTION_DECLARED_AVAILABLE,
CONNECTION_DECLARED_UNAVAILABLE,
InsightsSnapshot,
ProviderSnapshot,
SEVERITY_CRITICAL,
SEVERITY_INFO,
SEVERITY_UNPROVEN,
SEVERITY_WARN,
)
from webui.layout import render_page
_SEVERITY_CSS = {
SEVERITY_CRITICAL: "badge-blocked",
SEVERITY_WARN: "badge-health-degraded",
SEVERITY_INFO: "badge-health-ok",
SEVERITY_UNPROVEN: "badge-health-unproven",
}
_CONN_CSS = {
CONNECTION_DECLARED_AVAILABLE: "badge-health-ok",
CONNECTION_DECLARED_UNAVAILABLE: "badge-health-degraded",
"registry_unavailable": "badge-blocked",
}
def _badge(text: str, css: str) -> str:
return f'<span class="badge {css}">{escape(text)}</span>'
def _limits_card(lines: list[str], *, title: str) -> str:
items = "".join(f"<li>{escape(line)}</li>" for line in lines)
return f"""<div class="prompt-card">
<h3>{escape(title)}</h3>
<ul class="reasons">{items}</ul>
<p class="muted">Advisory surface only no review, merge, close, or provider
mutation is available here.</p>
</div>"""
def render_providers_page(snapshot: ProviderSnapshot) -> str:
"""Render the AI-provider connections page."""
if not snapshot.ok:
body = f"""<h2>AI provider connections</h2>
<p class="meta">Phase 4 read-only provider status (#650).</p>
<div class="health-card health-stale">
<strong>Provider registry unavailable:</strong>
{escape(snapshot.fetch_error or "registry could not be loaded")}.
No connection table is rendered an empty table would claim that no
providers are configured.
</div>
{_limits_card([
"connection_status reflects the worker registry declaration only",
"no API keys or credential material are loaded or rendered",
"live executable health is not probed on this surface",
], title="Interpretation limits")}
"""
return render_page(title="Providers", body_html=body)
rows = []
for provider in snapshot.providers:
models = (
", ".join(f"<code>{escape(m)}</code>" for m in provider.models)
if provider.models
else '<span class="muted">none declared</span>'
)
rows.append(
"<tr>"
f"<td><code>{escape(provider.provider_id)}</code></td>"
f"<td>{escape(provider.display_name)}</td>"
f"<td>{escape(provider.vendor)}</td>"
f"<td><code>{escape(provider.executable)}</code></td>"
f"<td>{_badge(provider.connection_status, _CONN_CSS.get(provider.connection_status, 'badge-health-skipped'))}</td>"
f"<td>{provider.worker_count} "
f"({provider.enabled_worker_count} enabled)</td>"
f"<td>{models}</td>"
"</tr>"
)
table = (
"".join(rows)
if rows
else '<tr><td colspan="7" class="muted">No providers declared in the registry.</td></tr>'
)
body = f"""<h2>AI provider connections</h2>
<p class="meta">Phase 4 read-only provider status (#650). Registry revision
<code>{escape(str(snapshot.registry_revision))}</code>.
Declared status only secrets never load.</p>
<div class="health-card">
<h3>Declared connections</h3>
<table class="registry">
<thead>
<tr>
<th>Provider</th><th>Name</th><th>Vendor</th><th>Executable</th>
<th>Connection</th><th>Workers</th><th>Models (declared)</th>
</tr>
</thead>
<tbody>{table}</tbody>
</table>
<p class="muted">{escape(snapshot.providers[0].probe_limit if snapshot.providers else "")}</p>
</div>
{_limits_card([
"connection_status reflects the worker registry declaration only",
"no API keys or credential material are loaded or rendered",
"live executable health is not probed on this surface (#800 owns that)",
], title="Interpretation limits")}
<p class="muted">Related: <a href="/insights">Operational insights</a> ·
<a href="/analytics">Analytics</a></p>
"""
return render_page(title="Providers", body_html=body)
def _evidence_list(insight) -> str:
items = "".join(
f"<li><code>{escape(ref.kind)}:{escape(ref.ref)}</code> — "
f"{escape(ref.detail)}</li>"
for ref in insight.evidence
)
return f'<ul class="reasons">{items}</ul>'
def render_insights_page(snapshot: InsightsSnapshot) -> str:
"""Render the operational insights page."""
if not snapshot.ok and not snapshot.insights:
body = f"""<h2>Operational insights</h2>
<p class="meta">Phase 4 evidence-backed insights (#650).</p>
<div class="health-card health-stale">
<strong>Insights unavailable:</strong>
{escape(snapshot.fetch_error or "no sources ran")}.
</div>
{_limits_card([
"insights are advisory only and never authorize merge, review, or close",
"an insight without evidence refs is refused rather than emitted",
], title="Interpretation limits")}
"""
return render_page(title="Insights", body_html=body)
source_bits = []
if snapshot.sources_used:
source_bits.append(
"sources used: " + ", ".join(f"<code>{escape(s)}</code>" for s in snapshot.sources_used)
)
if snapshot.sources_unavailable:
missing = "; ".join(
f"{escape(item.get('source', '?'))}: {escape(item.get('reason', ''))}"
for item in snapshot.sources_unavailable
)
source_bits.append(f"sources unavailable: {missing}")
cards = []
for insight in snapshot.insights:
cards.append(
f"""<div class="prompt-card">
<h3>{_badge(insight.severity, _SEVERITY_CSS.get(insight.severity, "badge-health-skipped"))}
{escape(insight.title)}</h3>
<p class="meta"><code>{escape(insight.kind)}</code> · confidence
<code>{escape(insight.confidence)}</code> ·
{_badge("advisory only", "badge-health-skipped")} ·
{_badge("no action claimed", "badge-health-ok")}</p>
<p>{escape(insight.summary)}</p>
<h4>Evidence</h4>
{_evidence_list(insight)}
</div>"""
)
if not cards:
cards.append(
'<div class="prompt-card"><p class="muted">No insights met the '
"evidence threshold in the loaded sources. That is not a claim "
"that the fleet is healthy — only that no qualifying pattern was "
"found.</p></div>"
)
body = f"""<h2>Operational insights</h2>
<p class="meta">Phase 4 evidence-backed insights (#650). Derived only from
durable console evidence; never invents policy or completes workflow actions.</p>
<p class="muted">{" · ".join(source_bits) if source_bits else ""}</p>
{"".join(cards)}
{_limits_card([
"insights are advisory only and never authorize merge, review, or close",
"an insight without evidence refs is refused rather than emitted",
"a missing source is listed as unavailable, not as an empty success",
"insights never claim a workflow action completed",
], title="Interpretation limits")}
<p class="muted">Related: <a href="/providers">Provider connections</a> ·
<a href="/traffic">Traffic</a> · <a href="/system-health">System health</a> ·
<a href="/analytics">Analytics</a></p>
"""
return render_page(title="Insights", body_html=body)
+14 -6
View File
@@ -4,11 +4,12 @@ Single source of truth for the console navigation so ``webui/layout.py`` and
the ``webui/app.py`` route table stay aligned with epic #631. Read-only: every
destination is a GET view or a Phase 1 placeholder. No mutation links.
Nav groups follow the #631 information architecture: Health, Traffic,
Nav groups follow the #631 Phase 1 information architecture: Health, Traffic,
Runtime/Sessions, Projects, Inventory, Timeline, Policy (placeholder), and
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.
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.
"""
from __future__ import annotations
@@ -45,9 +46,12 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
NavItem("/queue", "Queue"),
NavItem("/leases", "Leases"),
NavItem("/actions", "Actions"),
NavItem("/notifications", "Notifications"),
NavItem("/requests", "Requests"),
)),
NavGroup("Runtime/Sessions", (
NavItem("/runtime", "Runtime health"),
NavItem("/runtime/restart", "Restart status"),
NavItem("/sessions", "Sessions"),
)),
NavGroup("Projects", (
@@ -68,8 +72,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
NavItem("/prompts", "Prompts"),
)),
NavGroup("Insights", (
NavItem("/insights", "Insights"),
NavItem("/providers", "Providers"),
NavItem("/insights", "Insights", "stub"),
NavItem("/analytics", "Analytics"),
NavItem("/audit", "Audit"),
)),
@@ -93,6 +96,11 @@ 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.",
),
}
+158
View File
@@ -0,0 +1,158 @@
"""HTML rendering for Phase 3 Notifications and Human-Attention Console (#648)."""
from __future__ import annotations
from html import escape
from typing import Sequence
from webui.layout import render_page
from webui.notifications import (
ATTENTION_HUMAN_REQUIRED,
ATTENTION_OPERATOR,
ATTENTION_ROUTINE,
NotificationItem,
NotificationSnapshot,
)
def _render_attention_badge(attention_class: str) -> str:
cls = "badge"
if attention_class == ATTENTION_HUMAN_REQUIRED:
cls += " badge-blocked"
elif attention_class == ATTENTION_OPERATOR:
cls += " badge-claimed"
else:
cls += " muted"
return f'<span class="{cls}">{escape(attention_class)}</span>'
def _render_notification_row(item: NotificationItem) -> str:
category_label = escape(item.category.upper())
id_str = escape(item.id)
title_str = escape(item.title)
summary_str = escape(item.summary)
att_badge = _render_attention_badge(item.attention_class)
work_item_html = ""
if item.work_number and item.work_kind:
kind_label = escape(item.work_kind.upper())
num_str = f"#{item.work_number}"
link = item.deep_link or "#"
work_item_html = f'<a href="{escape(link)}"><code>{kind_label} {num_str}</code></a>'
requires_human_label = (
'<span class="badge badge-blocked" style="font-size:0.75rem;">HUMAN REQUIRED</span>'
if item.requires_human
else ""
)
return f"""<tr>
<td><code>{category_label}</code><br><span class="muted" style="font-size:0.75rem;">{id_str}</span></td>
<td>
<div><strong>{title_str}</strong> {att_badge} {requires_human_label}</div>
<div class="muted" style="font-size:0.85rem; margin-top:0.25rem;">{summary_str}</div>
</td>
<td>{work_item_html}</td>
<td><span class="muted" style="font-size:0.8rem;">{escape(item.created_at[:19])}</span></td>
</tr>"""
def _render_notifications_table(items: Sequence[NotificationItem], empty_message: str) -> str:
if not items:
return f'<p class="muted" style="padding:1rem 0;">{escape(empty_message)}</p>'
rows = "".join(_render_notification_row(item) for item in items)
return f"""<table class="registry">
<thead>
<tr>
<th style="width: 18%;">Category & ID</th>
<th style="width: 52%;">Title & Attention Summary</th>
<th style="width: 15%;">Work Item</th>
<th style="width: 15%;">Time</th>
</tr>
</thead>
<tbody>
{rows}
</tbody>
</table>"""
def render_notifications_page(
snapshot: NotificationSnapshot,
*,
filter_class: str = "inbox",
filter_project: str | None = None,
) -> str:
"""Render the notifications and attention inbox page."""
title = "Notifications & Attention Inbox"
err_html = ""
if snapshot.fetch_error:
err_html = f'<div class="stub" style="border-color:#e53e3e; background:#fff5f5; color:#c53030; margin-bottom:1rem;"><p><strong>Fetch Warning:</strong> {escape(snapshot.fetch_error)}</p></div>'
# Determine items to render based on filter_class
if filter_class == ATTENTION_HUMAN_REQUIRED:
display_items = snapshot.human_required_items
active_tab_title = "Human-Required Escalations"
elif filter_class == ATTENTION_OPERATOR:
display_items = snapshot.operator_items
active_tab_title = "Operator Inbox Items"
elif filter_class == ATTENTION_ROUTINE:
display_items = snapshot.routine_items
active_tab_title = "Routine Workflow Transitions"
elif filter_class == "all":
display_items = snapshot.items
active_tab_title = "All Events (including Routine)"
else: # "inbox" default
display_items = snapshot.inbox_items
active_tab_title = "Attention Inbox (Human + Operator)"
hr_cls = "badge-blocked" if snapshot.human_required_count > 0 else "muted"
op_cls = "badge-claimed" if snapshot.operator_count > 0 else "muted"
metrics_html = f"""<div style="display:flex; gap:1rem; margin-bottom:1.5rem;">
<div class="health-card" style="flex:1;">
<span class="muted" style="font-size:0.85rem;">Human Required</span>
<h2 style="margin:0.2rem 0;"><span class="badge {hr_cls}" style="font-size:1.4rem;">{snapshot.human_required_count}</span></h2>
<p class="muted" style="font-size:0.8rem; margin:0;">Critical escalation boundary</p>
</div>
<div class="health-card" style="flex:1;">
<span class="muted" style="font-size:0.85rem;">Operator Inbox</span>
<h2 style="margin:0.2rem 0;"><span class="badge {op_cls}" style="font-size:1.4rem;">{snapshot.operator_count}</span></h2>
<p class="muted" style="font-size:0.8rem; margin:0;">Operational items needing review</p>
</div>
<div class="health-card" style="flex:1;">
<span class="muted" style="font-size:0.85rem;">Routine Transitions</span>
<h2 style="margin:0.2rem 0;"><span class="badge muted" style="font-size:1.4rem;">{snapshot.routine_count}</span></h2>
<p class="muted" style="font-size:0.8rem; margin:0;">Background transitions (filtered)</p>
</div>
</div>"""
# Filter navigation links
def _tab_link(target_class: str, label: str) -> str:
is_active = (filter_class == target_class)
style = "font-weight:bold; border-bottom:2px solid currentColor;" if is_active else "color:#4a5568;"
return f'<a href="/notifications?attention_class={target_class}" style="margin-right:1.25rem; text-decoration:none; padding-bottom:0.25rem; {style}">{label}</a>'
tabs_html = f"""<div style="margin-bottom:1.25rem; border-bottom:1px solid #e2e8f0; padding-bottom:0.5rem;">
{_tab_link("inbox", f"Attention Inbox ({snapshot.human_required_count + snapshot.operator_count})")}
{_tab_link("human-required", f"Human Required ({snapshot.human_required_count})")}
{_tab_link("operator", f"Operator ({snapshot.operator_count})")}
{_tab_link("routine", f"Routine ({snapshot.routine_count})")}
{_tab_link("all", f"All Events ({snapshot.total_count})")}
</div>"""
table_html = _render_notifications_table(
display_items,
f"No items match attention filter '{filter_class}'.",
)
body = f"""<h2>{escape(title)}</h2>
<p class="muted">Phase 3 console surface for human-attention routing (#648). Routine workflow transitions are filtered by default to eliminate notification fatigue.</p>
{err_html}
{metrics_html}
{tabs_html}
<h3>{escape(active_tab_title)}</h3>
{table_html}"""
return render_page(title=title, body_html=body)
+486
View File
@@ -0,0 +1,486 @@
"""Notifications and human-attention routing module for Phase 3 web console (#648).
Defines attention classes, event classification rules, and inbox aggregation so
operators receive direct alerts only for human-required escalation boundaries
(#628) while routine workflow transitions remain available for pull-based review.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Callable
from webui import console_redaction
from webui.project_registry import load_registry
from webui.queue_loader import QueueSnapshot, load_queue_snapshot
from webui.lease_loader import LeaseSnapshot, load_lease_snapshot
from webui.system_health import SystemHealthSnapshot, load_system_health
# Attention class definitions (#628, #648)
ATTENTION_ROUTINE = "routine"
ATTENTION_OPERATOR = "operator"
ATTENTION_HUMAN_REQUIRED = "human-required"
ATTENTION_CLASSES = (
ATTENTION_ROUTINE,
ATTENTION_OPERATOR,
ATTENTION_HUMAN_REQUIRED,
)
# Notification categories
CATEGORY_AUTH = "auth"
CATEGORY_BLOCKER = "blocker"
CATEGORY_LEASE = "lease"
CATEGORY_VALIDATION = "validation"
CATEGORY_WORKFLOW = "workflow"
CATEGORY_SYSTEM = "system"
CATEGORIES = (
CATEGORY_AUTH,
CATEGORY_BLOCKER,
CATEGORY_LEASE,
CATEGORY_VALIDATION,
CATEGORY_WORKFLOW,
CATEGORY_SYSTEM,
)
@dataclass(frozen=True)
class NotificationItem:
"""A single notification or inbox event."""
id: str
attention_class: str # "routine", "operator", "human-required"
category: str # "auth", "blocker", "lease", "validation", etc.
title: str
summary: str
work_kind: str | None # "issue", "pr", "session", "system"
work_number: int | None
project_id: str
repo_label: str
created_at: str
deep_link: str | None = None
requires_human: bool = False
extra: dict[str, Any] = field(default_factory=dict)
def as_dict(self) -> dict[str, Any]:
return {
"id": self.id,
"attention_class": self.attention_class,
"category": self.category,
"title": self.title,
"summary": console_redaction.redact_text(self.summary),
"work_kind": self.work_kind,
"work_number": self.work_number,
"project_id": self.project_id,
"repo_label": self.repo_label,
"created_at": self.created_at,
"deep_link": self.deep_link,
"requires_human": self.requires_human,
"extra": self.extra,
}
@dataclass(frozen=True)
class NotificationSnapshot:
"""Snapshot of notifications and attention inbox state."""
project_id: str
repo_label: str
items: tuple[NotificationItem, ...]
human_required_count: int
operator_count: int
routine_count: int
total_count: int
fetch_error: str | None = None
@property
def inbox_items(self) -> tuple[NotificationItem, ...]:
"""Items requiring operator or human attention (excluding routine)."""
return tuple(
item
for item in self.items
if item.attention_class in {ATTENTION_OPERATOR, ATTENTION_HUMAN_REQUIRED}
)
@property
def human_required_items(self) -> tuple[NotificationItem, ...]:
return tuple(
item for item in self.items if item.attention_class == ATTENTION_HUMAN_REQUIRED
)
@property
def operator_items(self) -> tuple[NotificationItem, ...]:
return tuple(
item for item in self.items if item.attention_class == ATTENTION_OPERATOR
)
@property
def routine_items(self) -> tuple[NotificationItem, ...]:
return tuple(
item for item in self.items if item.attention_class == ATTENTION_ROUTINE
)
def as_dict(self) -> dict[str, Any]:
return {
"project_id": self.project_id,
"repo_label": self.repo_label,
"human_required_count": self.human_required_count,
"operator_count": self.operator_count,
"routine_count": self.routine_count,
"total_count": self.total_count,
"fetch_error": self.fetch_error,
"inbox_items": [item.as_dict() for item in self.inbox_items],
"all_items": [item.as_dict() for item in self.items],
}
def classify_attention_event(
category: str,
title: str,
summary: str,
*,
is_hard_stop: bool = False,
is_auth_failure: bool = False,
is_irrecoverable: bool = False,
is_decision_lock: bool = False,
is_validation_failure: bool = False,
is_stale: bool = False,
is_blocker: bool = False,
) -> tuple[str, bool]:
"""Classify an event into an attention class and human requirement flag.
Rules (#628, #648):
1. Critical boundaries (hard stop, auth failure, irrecoverable state,
decision lock, validation failure) -> ATTENTION_HUMAN_REQUIRED (requires_human=True).
2. Operational queues (blocker, stale lease, unassigned ready work, queue collision)
-> ATTENTION_OPERATOR (requires_human=False).
3. Routine state transitions (clean progression, healthy heartbeats) -> ATTENTION_ROUTINE (requires_human=False).
Classification uses structured flags and category only. Human-authored
``title`` / ``summary`` text is never substring-matched for escalation
(PR #905 review B1) — callers that need text signals must set flags from
machine-generated status/detail fields before calling this function.
"""
del title, summary # kept for API stability; never used for classification
if (
is_hard_stop
or is_auth_failure
or is_irrecoverable
or is_decision_lock
or is_validation_failure
or category in {CATEGORY_AUTH, CATEGORY_VALIDATION}
):
return ATTENTION_HUMAN_REQUIRED, True
if is_stale or is_blocker or category in {CATEGORY_BLOCKER, CATEGORY_LEASE}:
return ATTENTION_OPERATOR, False
return ATTENTION_ROUTINE, False
def load_notifications_snapshot(
project_id: str | None = None,
*,
load_queue: Callable[..., QueueSnapshot] | None = None,
load_leases: Callable[..., LeaseSnapshot] | None = None,
load_health: Callable[..., SystemHealthSnapshot] | None = None,
) -> NotificationSnapshot:
"""Load and classify attention notifications across queue, leases, and system health."""
registry = load_registry()
project = None
if project_id:
for entry in registry.projects:
if entry.id == project_id:
project = entry
break
else:
project = registry.projects[0] if registry.projects else None
if project is None:
return NotificationSnapshot(
project_id=project_id or "",
repo_label="",
items=(),
human_required_count=0,
operator_count=0,
routine_count=0,
total_count=0,
fetch_error="project not found in registry",
)
queue_loader_fn = load_queue or load_queue_snapshot
lease_loader_fn = load_leases or load_lease_snapshot
health_loader_fn = load_health or load_system_health
try:
queue_snap = queue_loader_fn(project.id)
except TypeError:
queue_snap = queue_loader_fn(project_id=project.id)
try:
lease_snap = lease_loader_fn(project_id=project.id)
except TypeError:
lease_snap = lease_loader_fn(project.id)
try:
health_snap = health_loader_fn(project_id=project.id)
except TypeError:
try:
health_snap = health_loader_fn(project.id)
except TypeError:
health_snap = health_loader_fn()
items: list[NotificationItem] = []
now_iso = datetime.now(timezone.utc).isoformat()
# 1. System health alerts (highest priority)
for err_idx, probe_err in enumerate(getattr(health_snap, "probe_errors", ())):
att_cls, req_human = classify_attention_event(
CATEGORY_SYSTEM,
"System Health Probe Error",
probe_err,
is_blocker=True,
)
items.append(
NotificationItem(
id=f"notif-sys-err-{project.id}-{err_idx}",
attention_class=att_cls,
category=CATEGORY_SYSTEM,
title="System Health Error",
summary=f"System health error: {probe_err}",
work_kind="system",
work_number=None,
project_id=project.id,
repo_label=f"{project.gitea_owner}/{project.repo_name}",
created_at=now_iso,
deep_link="/system",
requires_human=req_human,
)
)
for probe in getattr(health_snap, "dependencies", ()):
if probe.status not in ("ok", "healthy"):
att_cls, req_human = classify_attention_event(
CATEGORY_SYSTEM,
f"Probe Failure: {probe.name}",
probe.detail or probe.status,
is_hard_stop=("stop" in probe.status or "fatal" in probe.status),
is_auth_failure=("auth" in probe.name.lower() or "unauthorized" in probe.status.lower()),
is_blocker=True,
)
items.append(
NotificationItem(
id=f"notif-probe-{probe.name}",
attention_class=att_cls,
category=CATEGORY_AUTH if "auth" in probe.name.lower() else CATEGORY_SYSTEM,
title=f"Health Probe Alert: {probe.name}",
summary=f"Probe '{probe.name}' reported status '{probe.status}': {probe.detail}",
work_kind="system",
work_number=None,
project_id=project.id,
repo_label=f"{project.gitea_owner}/{project.repo_name}",
created_at=now_iso,
deep_link="/system",
requires_human=req_human,
)
)
# 2. Queue items (PRs and Issues)
for pr in queue_snap.prs:
if "blocked" in pr.badges:
att_cls, req_human = classify_attention_event(
CATEGORY_BLOCKER,
f"PR #{pr.number} Blocked",
f"PR #{pr.number} '{pr.title}' is blocked or has merge conflicts.",
is_blocker=True,
)
items.append(
NotificationItem(
id=f"notif-pr-block-{pr.number}",
attention_class=att_cls,
category=CATEGORY_BLOCKER,
title=f"Blocked PR #{pr.number}",
summary=f"PR #{pr.number} ({pr.title}) requires merge conflict resolution.",
work_kind="pr",
work_number=pr.number,
project_id=project.id,
repo_label=f"{project.gitea_owner}/{project.repo_name}",
created_at=now_iso,
deep_link=f"/traffic",
requires_human=req_human,
)
)
elif "stale" in pr.badges:
att_cls, req_human = classify_attention_event(
CATEGORY_WORKFLOW,
f"PR #{pr.number} Stale",
f"PR #{pr.number} '{pr.title}' has had no activity for over 14 days.",
is_stale=True,
)
items.append(
NotificationItem(
id=f"notif-pr-stale-{pr.number}",
attention_class=att_cls,
category=CATEGORY_WORKFLOW,
title=f"Stale PR #{pr.number}",
summary=f"PR #{pr.number} ({pr.title}) is stale.",
work_kind="pr",
work_number=pr.number,
project_id=project.id,
repo_label=f"{project.gitea_owner}/{project.repo_name}",
created_at=now_iso,
deep_link=f"/queue",
requires_human=req_human,
)
)
else:
# Routine PR transition
att_cls, req_human = classify_attention_event(
CATEGORY_WORKFLOW,
f"PR #{pr.number} Active",
f"PR #{pr.number} '{pr.title}' is in routine state {', '.join(pr.badges)}.",
)
items.append(
NotificationItem(
id=f"notif-pr-routine-{pr.number}",
attention_class=att_cls,
category=CATEGORY_WORKFLOW,
title=f"Routine PR #{pr.number}",
summary=f"PR #{pr.number} ({pr.title}) state: {', '.join(pr.badges)}.",
work_kind="pr",
work_number=pr.number,
project_id=project.id,
repo_label=f"{project.gitea_owner}/{project.repo_name}",
created_at=now_iso,
deep_link=f"/queue",
requires_human=req_human,
)
)
for issue in queue_snap.issues:
if "duplicate" in issue.badges:
att_cls, req_human = classify_attention_event(
CATEGORY_BLOCKER,
f"Issue #{issue.number} Duplicate PRs",
f"Issue #{issue.number} has multiple linked PRs.",
is_blocker=True,
)
items.append(
NotificationItem(
id=f"notif-issue-dup-{issue.number}",
attention_class=att_cls,
category=CATEGORY_BLOCKER,
title=f"Duplicate PRs on Issue #{issue.number}",
summary=f"Issue #{issue.number} ({issue.title}) linked to multiple PRs.",
work_kind="issue",
work_number=issue.number,
project_id=project.id,
repo_label=f"{project.gitea_owner}/{project.repo_name}",
created_at=now_iso,
deep_link=f"/traffic",
requires_human=req_human,
)
)
elif "claimed" in issue.badges or "in-review" in issue.badges:
att_cls, req_human = classify_attention_event(
CATEGORY_WORKFLOW,
f"Issue #{issue.number} Active",
f"Issue #{issue.number} '{issue.title}' in state {', '.join(issue.badges)}.",
)
items.append(
NotificationItem(
id=f"notif-issue-routine-{issue.number}",
attention_class=att_cls,
category=CATEGORY_WORKFLOW,
title=f"Routine Issue #{issue.number}",
summary=f"Issue #{issue.number} ({issue.title}) state: {', '.join(issue.badges)}.",
work_kind="issue",
work_number=issue.number,
project_id=project.id,
repo_label=f"{project.gitea_owner}/{project.repo_name}",
created_at=now_iso,
deep_link=f"/queue",
requires_human=req_human,
)
)
# 3. Leases / Collisions
for lease in lease_snap.reviewer_leases:
if lease.get("is_expired") or lease.get("status") == "expired":
pr_num = lease.get("pr_number") or lease.get("work_item_number")
att_cls, req_human = classify_attention_event(
CATEGORY_LEASE,
f"Reviewer Lease Expired for PR #{pr_num}",
f"Reviewer lease for PR #{pr_num} has expired.",
is_stale=True,
)
items.append(
NotificationItem(
id=f"notif-lease-exp-pr-{pr_num}",
attention_class=att_cls,
category=CATEGORY_LEASE,
title=f"Expired Reviewer Lease (PR #{pr_num})",
summary=f"Reviewer lease for PR #{pr_num} expired.",
work_kind="pr",
work_number=pr_num,
project_id=project.id,
repo_label=f"{project.gitea_owner}/{project.repo_name}",
created_at=now_iso,
deep_link="/leases",
requires_human=req_human,
)
)
for col_idx, collision in enumerate(lease_snap.duplicate_prs):
att_cls, req_human = classify_attention_event(
CATEGORY_BLOCKER,
f"Duplicate PR Collision ({collision.kind})",
collision.message,
is_blocker=True,
)
issue_part = collision.issue_number if collision.issue_number is not None else "none"
kind_part = (collision.kind or "unknown").replace(" ", "-")
items.append(
NotificationItem(
id=f"notif-collision-{kind_part}-{issue_part}-{col_idx}",
attention_class=att_cls,
category=CATEGORY_BLOCKER,
title=f"Collision Alert ({collision.kind})",
summary=collision.message,
work_kind="issue" if collision.issue_number else "pr",
work_number=collision.issue_number,
project_id=project.id,
repo_label=f"{project.gitea_owner}/{project.repo_name}",
created_at=now_iso,
deep_link="/leases",
requires_human=req_human,
)
)
human_req_count = sum(1 for i in items if i.attention_class == ATTENTION_HUMAN_REQUIRED)
operator_count = sum(1 for i in items if i.attention_class == ATTENTION_OPERATOR)
routine_count = sum(1 for i in items if i.attention_class == ATTENTION_ROUTINE)
# Fetch errors are transport/load failures only — not probe results that
# already surface as first-class notification items (PR #905 review B3).
fetch_err = queue_snap.fetch_error or lease_snap.fetch_error
if isinstance(fetch_err, (tuple, list)):
fetch_err = "; ".join(fetch_err) if fetch_err else None
return NotificationSnapshot(
project_id=project.id,
repo_label=f"{project.gitea_owner}/{project.repo_name}",
items=tuple(items),
human_required_count=human_req_count,
operator_count=operator_count,
routine_count=routine_count,
total_count=len(items),
fetch_error=fetch_err,
)
def snapshot_to_dict(snapshot: NotificationSnapshot) -> dict[str, Any]:
"""JSON-serializable export for /api/v1/notifications."""
return snapshot.as_dict()
File diff suppressed because it is too large Load Diff
+164
View File
@@ -0,0 +1,164 @@
"""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)
+579
View File
@@ -0,0 +1,579 @@
"""Read-only restart status, impact preview, and approval state (#667).
Phase 1 of the console restart surface. It *consumes* the #655 coordinator
substrate and renders it; it never restarts, reloads, drains, approves, or kills
anything. There is no apply path in this module, so there is no execution gate
here to arm incorrectly the only writes the console could perform are the ones
it does not implement.
Sources, each independently fail-soft and each reported with its own
:class:`SourceStatus`:
* :mod:`restart_coordinator` restart-class policy matrix (#663) and the
blast-radius impact report (#658).
* :mod:`drain_proof` drain checklist and gate verdict (#661), verified
read-only against a caller-supplied proof.
* :mod:`post_restart_reconcile` post-restart completion proof (#662).
* :mod:`webui.console_authz` role authorization for the approval controls
(#633).
Three rules this module holds itself to, because a status surface that lies is
worse than one that is absent:
**A source that could not be read is reported unavailable, never green.** No
default, placeholder, or self-comparison is substituted for a reading that
failed. An unreadable control-plane DB yields ``inventory_complete=False``,
which the coordinator itself turns into a fail-closed verdict.
**Authorization is asked the way execution would ask it.** Every authorization
probe passes ``for_execution=True``, so the console reports whether the action
could actually run rather than the weaker "this principal is the right role".
While the console is in Phase 1 that answer is ``phase_not_active`` for every
phase-2 action, and the surface says so plainly instead of showing an allow.
**The database is opened read-only.** ``ControlPlaneDB()`` creates directories
and runs migrations on construction, which is a write; this module opens the
sqlite file with ``mode=ro`` exactly as :mod:`webui.inventory` does, and treats
a missing file as missing authority rather than an empty inventory.
"""
from __future__ import annotations
import os
import sqlite3
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Callable, Mapping
import control_plane_db
import drain_proof
import restart_coordinator
from webui import console_authz
from webui.inventory import redact_path, scrub
# --- Source status ----------------------------------------------------------
STATUS_OK = "ok"
STATUS_UNAVAILABLE = "unavailable"
#: Console actions whose authorization state this surface reports. Both are
#: pre-existing #642 actions; this module adds no new console action because it
#: performs no console action.
REPORTED_ACTIONS: tuple[str, ...] = (
"system.restart_namespace",
"system.reload_namespace",
)
#: The break-glass workflow (#664) is not consumed here. It is declared so the
#: surface is honest about the gap rather than silently omitting a governance
#: path the operator has been told exists.
BREAK_GLASS_ISSUE = 664
BREAK_GLASS_PENDING_REASON = (
"The break-glass workflow (#664) is not yet available on this branch's "
"base; no break-glass control is offered and none is implied."
)
@dataclass(frozen=True)
class SourceStatus:
"""Whether one backing source could be read, and why not when it could not."""
name: str
status: str
detail: str = ""
@property
def available(self) -> bool:
return self.status == STATUS_OK
def as_dict(self) -> dict[str, Any]:
return {
"name": self.name,
"status": self.status,
"available": self.available,
"detail": self.detail,
}
@dataclass(frozen=True)
class RestartClassView:
"""One row of the #663 restart-class matrix, scoped to the viewer's role."""
restart_class: str
required_permission: str
expected_blast_radius: str
drain_requirement: str
full_drain_required: bool
approval_requirement: str
request_roles: tuple[str, ...]
execution_roles: tuple[str, ...]
viewer_may_request: bool
viewer_may_execute: bool
def as_dict(self) -> dict[str, Any]:
return {
"restart_class": self.restart_class,
"required_permission": self.required_permission,
"expected_blast_radius": self.expected_blast_radius,
"drain_requirement": self.drain_requirement,
"full_drain_required": self.full_drain_required,
"approval_requirement": self.approval_requirement,
"request_roles": list(self.request_roles),
"execution_roles": list(self.execution_roles),
"viewer_may_request": self.viewer_may_request,
"viewer_may_execute": self.viewer_may_execute,
}
@dataclass(frozen=True)
class ActionAuthorization:
"""Authorization state for one console action, asked as execution would."""
action_id: str
summary: str
required_role: str
allowed: bool
execution_enabled: bool
reason_code: str
detail: str
def as_dict(self) -> dict[str, Any]:
return {
"action_id": self.action_id,
"summary": self.summary,
"required_role": self.required_role,
"allowed": self.allowed,
"execution_enabled": self.execution_enabled,
"reason_code": self.reason_code,
"detail": self.detail,
}
@dataclass(frozen=True)
class BreakGlassSurface:
"""Declared-but-unavailable break-glass panel (#664 is not on this base)."""
available: bool
issue: int
reason: str
viewer_is_privileged: bool
def as_dict(self) -> dict[str, Any]:
return {
"available": self.available,
"issue": self.issue,
"reason": self.reason,
"viewer_is_privileged": self.viewer_is_privileged,
}
@dataclass(frozen=True)
class RestartConsoleSnapshot:
"""Everything the read-only restart console renders."""
generated_at: str
viewer_role: str
viewer_authenticated: bool
read_only: bool
impact: dict[str, Any] | None
impact_source: SourceStatus
drain: dict[str, Any] | None
drain_source: SourceStatus
reconcile: dict[str, Any] | None
reconcile_source: SourceStatus
restart_classes: tuple[RestartClassView, ...]
authorizations: tuple[ActionAuthorization, ...]
break_glass: BreakGlassSurface
notes: tuple[str, ...] = field(default_factory=tuple)
def as_dict(self) -> dict[str, Any]:
return {
"generated_at": self.generated_at,
"viewer_role": self.viewer_role,
"viewer_authenticated": self.viewer_authenticated,
"read_only": self.read_only,
"impact": self.impact,
"impact_source": self.impact_source.as_dict(),
"drain": self.drain,
"drain_source": self.drain_source.as_dict(),
"reconcile": self.reconcile,
"reconcile_source": self.reconcile_source.as_dict(),
"restart_classes": [c.as_dict() for c in self.restart_classes],
"authorizations": [a.as_dict() for a in self.authorizations],
"break_glass": self.break_glass.as_dict(),
"notes": list(self.notes),
"links": {
"issue": 667,
"extends": 642,
"umbrella": 655,
"coordinator": 658,
"drain_proof": 661,
"reconcile": 662,
"restart_classes": 663,
"break_glass": BREAK_GLASS_ISSUE,
"vision": 652,
"roadmap": 653,
},
}
def _utc_now() -> datetime:
return datetime.now(timezone.utc)
# --- Control-plane inventory (read-only) ------------------------------------
def read_control_plane_inventory(
*,
db_path: str | None = None,
limit: int = 200,
) -> dict[str, Any]:
"""Read sessions and leases for an impact evaluation, read-only.
Returns the inventory mapping
:func:`restart_coordinator.evaluate_restart_impact` expects.
``inventory_complete`` is True only when every read succeeded, so a partial
read denies rather than under-reporting the blast radius.
The database is never created, migrated, or written: a missing file means
the console has no session authority, which is not the same as there being
no sessions.
"""
path = (db_path or control_plane_db.default_db_path() or "").strip()
incomplete: list[str] = []
def _incomplete(reason: str) -> dict[str, Any]:
return {
"sessions": [],
"leases": [],
"terminal_lock": None,
"prior_recovery_attempts": [],
"inventory_complete": False,
"incomplete_reasons": [reason],
}
if not path:
return _incomplete("control-plane database path is not configured")
if not os.path.exists(path):
return _incomplete(
f"control-plane database not present at {redact_path(path)}; "
"no session or lease authority available"
)
try:
conn = sqlite3.connect(f"file:{path}?mode=ro", uri=True, timeout=5)
conn.row_factory = sqlite3.Row
except sqlite3.Error as exc:
return _incomplete(f"control-plane database could not be opened: {exc}")
sessions: list[dict[str, Any]] = []
leases: list[dict[str, Any]] = []
capped = max(1, int(limit))
try:
tables = {
str(row[0])
for row in conn.execute(
"SELECT name FROM sqlite_master WHERE type = 'table'"
).fetchall()
}
if "sessions" not in tables:
incomplete.append("control-plane database has no sessions table")
else:
sessions = [
dict(row)
for row in conn.execute(
"SELECT session_id, role, profile, pid, status,"
" last_heartbeat_at FROM sessions"
" WHERE status = 'active'"
" ORDER BY last_heartbeat_at DESC LIMIT ?",
(capped,),
).fetchall()
]
if "leases" not in tables:
incomplete.append("control-plane database has no leases table")
elif "work_items" not in tables:
incomplete.append(
"control-plane database has no work_items table; lease work "
"identity cannot be resolved"
)
else:
leases = [
dict(row)
for row in conn.execute(
"SELECT l.lease_id, l.session_id, l.role, l.phase,"
" l.status AS freshness, l.worktree_path,"
" w.kind AS work_kind, w.number AS work_number"
" FROM leases l"
" JOIN work_items w ON w.work_item_id = l.work_item_id"
" WHERE l.status = 'active'"
" ORDER BY l.expires_at DESC LIMIT ?",
(capped,),
).fetchall()
]
except sqlite3.Error as exc:
return _incomplete(f"control-plane database read failed: {exc}")
finally:
conn.close()
return {
"sessions": sessions,
"leases": leases,
"terminal_lock": None,
"prior_recovery_attempts": [],
"inventory_complete": not incomplete,
"incomplete_reasons": incomplete,
}
# --- Composition ------------------------------------------------------------
def build_restart_class_views(viewer_role: str | None) -> tuple[RestartClassView, ...]:
"""Render the #663 class matrix, marking what this viewer may request."""
normalized = str(viewer_role or "").strip().lower()
views: list[RestartClassView] = []
for policy in restart_coordinator.RESTART_CLASS_POLICIES.values():
views.append(
RestartClassView(
restart_class=policy.restart_class.value,
required_permission=policy.required_permission,
expected_blast_radius=policy.expected_blast_radius,
drain_requirement=policy.drain_requirement,
full_drain_required=policy.full_drain_required,
approval_requirement=policy.approval_requirement,
request_roles=tuple(policy.request_roles),
execution_roles=tuple(policy.execution_roles),
viewer_may_request=normalized in policy.request_roles,
viewer_may_execute=normalized in policy.execution_roles,
)
)
return tuple(views)
def build_action_authorizations(
principal: console_authz.Principal | None,
) -> tuple[ActionAuthorization, ...]:
"""Authorization state for the approval controls, asked as execution.
``for_execution=True`` is deliberate. Asking without it answers "is this
principal senior enough", which is not the question an operator looking at a
control needs answered; asking with it answers "would this run", and while
the console is in Phase 1 the honest answer is no.
"""
results: list[ActionAuthorization] = []
for action_id in REPORTED_ACTIONS:
action = console_authz.get_action(action_id)
decision = console_authz.authorize(action_id, principal, for_execution=True)
results.append(
ActionAuthorization(
action_id=action_id,
summary=action.summary if action else "",
required_role=(
action.minimum_role if action else console_authz.OPERATOR
),
allowed=bool(decision.allowed),
execution_enabled=bool(decision.execution_enabled),
reason_code=str(decision.reason_code or ""),
detail=str(decision.detail or ""),
)
)
return tuple(results)
def viewer_is_privileged(principal: console_authz.Principal | None) -> bool:
"""True when the viewer holds at least the operator role."""
who = principal if principal is not None else console_authz.ANONYMOUS
if not who.authenticated:
return False
return who.rank >= console_authz.ROLE_ORDER.index(console_authz.OPERATOR)
def load_impact_report(
*,
principal: console_authz.Principal | None = None,
restart_class: str = restart_coordinator.RestartClass.FULL_MCP_RESTART.value,
db_path: str | None = None,
limit: int = 200,
read_inventory: Callable[..., Mapping[str, Any]] | None = None,
now: datetime | None = None,
) -> tuple[dict[str, Any] | None, SourceStatus]:
"""Evaluate the blast radius for *restart_class*, always dry-run."""
reader = read_inventory or read_control_plane_inventory
try:
inventory = dict(reader(db_path=db_path, limit=limit))
except Exception as exc: # noqa: BLE001
return None, SourceStatus(
"impact",
STATUS_UNAVAILABLE,
f"control-plane inventory failed: {type(exc).__name__}: {exc}",
)
who = principal if principal is not None else console_authz.ANONYMOUS
viewer_role = str(who.role or "").strip().lower()
try:
report = restart_coordinator.evaluate_restart_impact(
inventory,
now=now,
dry_run=True,
restart_class=restart_class,
requester_role=viewer_role,
requester_permissions=restart_coordinator.permissions_for_role(
viewer_role
),
)
except Exception as exc: # noqa: BLE001
return None, SourceStatus(
"impact",
STATUS_UNAVAILABLE,
f"impact evaluation failed: {type(exc).__name__}: {exc}",
)
payload = scrub(report.as_dict())
detail = ""
if not report.inventory_complete:
detail = "; ".join(report.incomplete_reasons) or "inventory incomplete"
return payload, SourceStatus("impact", STATUS_OK, detail)
def load_drain_status(
*,
proof: Mapping[str, Any] | None = None,
now: datetime | None = None,
expected_impact_fingerprint: str | None = None,
) -> tuple[dict[str, Any] | None, SourceStatus]:
"""Verify a supplied drain proof read-only and report the verdict.
No proof supplied is not a failure and not a pass: it is reported as the
absence of a proof, which is exactly what the #661 gate would deny on.
"""
if proof is None:
return None, SourceStatus(
"drain",
STATUS_UNAVAILABLE,
"no drain proof supplied; the #661 gate denies a restart without a "
"valid unexpired clean proof",
)
try:
verified = drain_proof.verify_drain_proof(
proof,
now=now,
expected_impact_fingerprint=expected_impact_fingerprint,
)
except Exception as exc: # noqa: BLE001
return None, SourceStatus(
"drain",
STATUS_UNAVAILABLE,
f"drain proof verification failed: {type(exc).__name__}: {exc}",
)
return scrub(verified.as_dict()), SourceStatus("drain", STATUS_OK)
def load_reconcile_status(
*,
load_proof: Callable[[], Any] | None = None,
) -> tuple[dict[str, Any] | None, SourceStatus]:
"""Report the most recent post-restart completion proof (#662)."""
if load_proof is None:
return None, SourceStatus(
"reconcile",
STATUS_UNAVAILABLE,
"no post-restart completion proof source is wired into this view",
)
try:
proof = load_proof()
except Exception as exc: # noqa: BLE001
return None, SourceStatus(
"reconcile",
STATUS_UNAVAILABLE,
f"reconcile proof unavailable: {type(exc).__name__}: {exc}",
)
if proof is None:
return None, SourceStatus(
"reconcile",
STATUS_UNAVAILABLE,
"no post-restart reconcile has been recorded",
)
payload = proof.as_dict() if hasattr(proof, "as_dict") else dict(proof)
return scrub(payload), SourceStatus("reconcile", STATUS_OK)
def load_restart_console_snapshot(
*,
principal: console_authz.Principal | None = None,
restart_class: str = restart_coordinator.RestartClass.FULL_MCP_RESTART.value,
db_path: str | None = None,
limit: int = 200,
drain_proof_payload: Mapping[str, Any] | None = None,
read_inventory: Callable[..., Mapping[str, Any]] | None = None,
load_reconcile_proof: Callable[[], Any] | None = None,
now: datetime | None = None,
) -> RestartConsoleSnapshot:
"""Compose the read-only restart console snapshot."""
who = principal if principal is not None else console_authz.ANONYMOUS
moment = now or _utc_now()
impact, impact_source = load_impact_report(
principal=who,
restart_class=restart_class,
db_path=db_path,
limit=limit,
read_inventory=read_inventory,
now=moment,
)
fingerprint = None
if impact is not None:
try:
fingerprint = drain_proof.impact_fingerprint(impact)
except Exception: # noqa: BLE001
fingerprint = None
drain, drain_source = load_drain_status(
proof=drain_proof_payload,
now=moment,
expected_impact_fingerprint=fingerprint,
)
reconcile, reconcile_source = load_reconcile_status(
load_proof=load_reconcile_proof
)
notes: list[str] = [
"This surface is read-only: it evaluates and displays, and performs no "
"restart, reload, drain, approval, or process action.",
]
if not impact_source.available:
notes.append(
"Impact preview unavailable — a restart decision must not be made "
"from this page while the blast radius is unknown."
)
return RestartConsoleSnapshot(
generated_at=moment.isoformat(),
viewer_role=str(who.role or "anonymous"),
viewer_authenticated=bool(who.authenticated),
read_only=True,
impact=impact,
impact_source=impact_source,
drain=drain,
drain_source=drain_source,
reconcile=reconcile,
reconcile_source=reconcile_source,
restart_classes=build_restart_class_views(who.role),
authorizations=build_action_authorizations(who),
break_glass=BreakGlassSurface(
available=False,
issue=BREAK_GLASS_ISSUE,
reason=BREAK_GLASS_PENDING_REASON,
viewer_is_privileged=viewer_is_privileged(who),
),
notes=tuple(notes),
)
+299
View File
@@ -0,0 +1,299 @@
"""HTML views for the read-only restart console (#667).
Every interpolated value passes through :func:`_esc`. Values that can carry a
filesystem path or free-form operator text additionally pass through
:func:`webui.inventory.scrub_text`, which redacts credential-shaped tokens
*inside* a string rather than only at its start.
The page renders state and never offers a control that would mutate anything:
the approval and break-glass panels report authorization and availability, and
there is no form, button, or endpoint behind them.
"""
from __future__ import annotations
import html
from webui.inventory import scrub_text
from webui.restart_console import RestartConsoleSnapshot, SourceStatus
def _esc(value: object) -> str:
"""Escape any value for HTML text or a quoted attribute."""
if value is None:
return ""
return html.escape(str(value), quote=True)
def _esc_text(value: object) -> str:
"""Escape free-form text after redacting secrets embedded inside it."""
if value is None:
return ""
return _esc(scrub_text(str(value)))
def _bool_badge(
value: bool, *, true_label: str = "yes", false_label: str = "no"
) -> str:
css = "badge-ok" if value else "badge-blocked"
label = true_label if value else false_label
return f'<span class="badge {css}">{_esc(label)}</span>'
def _source_badge(source: SourceStatus) -> str:
css = "badge-ok" if source.available else "badge-blocked"
badge = f'<span class="badge {css}">{_esc(source.status)}</span>'
if source.detail:
badge += f' <span class="muted">{_esc_text(source.detail)}</span>'
return badge
def _notes_block(snapshot: RestartConsoleSnapshot) -> str:
if not snapshot.notes:
return ""
items = "".join(f"<li>{_esc_text(note)}</li>" for note in snapshot.notes)
return f"<ul class='reasons'>{items}</ul>"
def _impact_section(snapshot: RestartConsoleSnapshot) -> str:
head = (
"<section class='health-card'>"
f"<h3>Impact preview {_source_badge(snapshot.impact_source)}</h3>"
)
impact = snapshot.impact
if impact is None:
return (
head
+ "<p class='muted'>No impact preview is available, so the blast "
"radius of a restart is unknown. Treat this as unsafe.</p></section>"
)
counts = impact.get("counts") or {}
verdict = str(impact.get("verdict") or "unknown")
verdict_css = "badge-ok" if verdict == "safe" else "badge-blocked"
rows = "".join(
f"<tr><th>{_esc(key.replace('_', ' '))}</th><td>{_esc(value)}</td></tr>"
for key, value in sorted(counts.items())
)
reasons = "".join(
f"<li>{_esc_text(reason)}</li>" for reason in (impact.get("reasons") or [])
)
incomplete = ""
if not impact.get("inventory_complete", False):
detail = "; ".join(str(r) for r in (impact.get("incomplete_reasons") or []))
incomplete = (
"<p class='error'><strong>Inventory incomplete:</strong> "
f"{_esc_text(detail or 'unspecified')}. The coordinator fails "
"closed on an incomplete inventory.</p>"
)
sessions = impact.get("affected_sessions") or []
session_rows = "".join(
"<tr>"
f"<td><code>{_esc(s.get('session_id'))}</code></td>"
f"<td>{_esc(s.get('role'))}</td>"
f"<td>{_esc(s.get('pid'))}</td>"
f"<td>{_bool_badge(bool(s.get('live')), true_label='live', false_label='idle')}</td>"
f"<td>{_bool_badge(not s.get('heartbeat_stale'), true_label='fresh', false_label='stale')}</td>"
"</tr>"
for s in sessions[:50]
)
session_table = (
"<h4>Sessions a restart would terminate</h4>"
"<div class='table-scroll'><table class='registry'><thead><tr>"
"<th>Session</th><th>Role</th><th>PID</th><th>State</th>"
"<th>Heartbeat</th></tr></thead><tbody>"
f"{session_rows}</tbody></table></div>"
if session_rows
else "<p class='muted'>No affected sessions reported.</p>"
)
truncated = (
f"<p class='muted'>Showing the first 50 of {_esc(len(sessions))} "
"affected sessions.</p>"
if len(sessions) > 50
else ""
)
return (
head
+ "<p class='health-headline'>Verdict "
f"<span class='badge {verdict_css}'>{_esc(verdict)}</span> · "
f"blast radius <code>{_esc(impact.get('blast_radius'))}</code> · "
f"class <code>{_esc(impact.get('restart_class'))}</code></p>"
+ incomplete
+ (f"<ul class='reasons'>{reasons}</ul>" if reasons else "")
+ (f"<table class='registry'><tbody>{rows}</tbody></table>" if rows else "")
+ session_table
+ truncated
+ "</section>"
)
def _drain_section(snapshot: RestartConsoleSnapshot) -> str:
head = (
"<section class='health-card'>"
f"<h3>Drain proof {_source_badge(snapshot.drain_source)}</h3>"
)
drain = snapshot.drain
if drain is None:
return (
head
+ "<p class='muted'>No drain proof has been presented to this view. "
"The #661 gate authorizes a restart only against a valid, unexpired, "
"clean proof, so the absence of one is a denial, not a pass.</p>"
"</section>"
)
reasons = "".join(
f"<li>{_esc_text(reason)}</li>" for reason in (drain.get("reasons") or [])
)
return (
head
+ "<table class='registry'><tbody>"
f"<tr><th>Valid</th><td>{_bool_badge(bool(drain.get('valid')))}</td></tr>"
f"<tr><th>Clean</th><td>{_bool_badge(bool(drain.get('clean')))}</td></tr>"
f"<tr><th>Expired</th><td>{_bool_badge(not drain.get('expired'), true_label='no', false_label='yes')}</td></tr>"
f"<tr><th>Tampered</th><td>{_bool_badge(not drain.get('tampered'), true_label='no', false_label='yes')}</td></tr>"
f"<tr><th>Proof id</th><td><code>{_esc(drain.get('proof_id'))}</code></td></tr>"
"</tbody></table>"
+ (f"<ul class='reasons'>{reasons}</ul>" if reasons else "")
+ "</section>"
)
def _reconcile_section(snapshot: RestartConsoleSnapshot) -> str:
head = (
"<section class='health-card'>"
f"<h3>Post-restart reconcile {_source_badge(snapshot.reconcile_source)}</h3>"
)
proof = snapshot.reconcile
if proof is None:
return (
head
+ "<p class='muted'>No post-restart completion proof is recorded. "
"Until one is, the last restart's recovery state is unproven.</p>"
"</section>"
)
items = "".join(
"<tr>"
f"<td>{_esc(item.get('dimension'))}</td>"
f"<td>{_esc(item.get('status'))}</td>"
f"<td>{_esc_text(item.get('summary'))}</td>"
f"<td>{_bool_badge(not item.get('follow_up_required'), true_label='no', false_label='yes')}</td>"
"</tr>"
for item in (proof.get("items") or [])
)
return (
head
+ "<p class='health-headline'>Status "
f"<code>{_esc(proof.get('overall_status'))}</code> · mode "
f"<code>{_esc(proof.get('mode'))}</code> · resolved "
f"{_esc(proof.get('resolved_count'))} · unresolved "
f"{_esc(proof.get('unresolved_count'))}</p>"
+ (
"<div class='table-scroll'><table class='registry'><thead><tr>"
"<th>Dimension</th><th>Status</th><th>Summary</th>"
"<th>Follow-up required</th></tr></thead><tbody>"
f"{items}</tbody></table></div>"
if items
else "<p class='muted'>No reconcile dimensions reported.</p>"
)
+ "</section>"
)
def _class_matrix_section(snapshot: RestartConsoleSnapshot) -> str:
rows = "".join(
"<tr>"
f"<td><code>{_esc(view.restart_class)}</code></td>"
f"<td><code>{_esc(view.required_permission)}</code></td>"
f"<td>{_esc(view.expected_blast_radius)}</td>"
f"<td>{_esc(view.drain_requirement)}</td>"
f"<td>{_esc(view.approval_requirement)}</td>"
f"<td>{_bool_badge(view.viewer_may_request)}</td>"
f"<td>{_bool_badge(view.viewer_may_execute)}</td>"
"</tr>"
for view in snapshot.restart_classes
)
return (
"<section class='health-card'>"
"<h3>Restart classes</h3>"
"<p class='muted'>The least-privilege matrix each restart request is "
"resolved against. &ldquo;You may request&rdquo; and &ldquo;you may "
"execute&rdquo; are computed for the current viewer role, not for a "
"generic operator.</p>"
"<div class='table-scroll'><table class='registry'><thead><tr>"
"<th>Class</th><th>Permission</th><th>Blast radius</th>"
"<th>Drain</th><th>Approval</th><th>You may request</th>"
"<th>You may execute</th></tr></thead><tbody>"
f"{rows}</tbody></table></div>"
"</section>"
)
def _approval_section(snapshot: RestartConsoleSnapshot) -> str:
rows = "".join(
"<tr>"
f"<td><code>{_esc(a.action_id)}</code></td>"
f"<td>{_esc(a.required_role)}</td>"
f"<td>{_bool_badge(a.allowed)}</td>"
f"<td>{_bool_badge(a.execution_enabled)}</td>"
f"<td><code>{_esc(a.reason_code)}</code></td>"
f"<td>{_esc_text(a.detail)}</td>"
"</tr>"
for a in snapshot.authorizations
)
return (
"<section class='health-card'>"
"<h3>Approval controls</h3>"
"<p class='muted'>Authorization is probed the way execution would probe "
"it, so &ldquo;execution enabled&rdquo; answers whether the action would "
"actually run — not merely whether this role outranks the requirement. "
"No control on this page performs the action.</p>"
"<div class='table-scroll'><table class='registry'><thead><tr>"
"<th>Action</th><th>Required role</th><th>Authorized</th>"
"<th>Execution enabled</th><th>Reason</th><th>Detail</th>"
"</tr></thead><tbody>"
f"{rows}</tbody></table></div>"
"</section>"
)
def _break_glass_section(snapshot: RestartConsoleSnapshot) -> str:
bg = snapshot.break_glass
if not bg.viewer_is_privileged:
return (
"<section class='health-card'>"
"<h3>Break-glass</h3>"
"<p class='muted'>Break-glass status is visible to operator-class "
"roles only. Your role does not carry that authority, so no "
"emergency surface is shown.</p>"
"</section>"
)
return (
"<section class='health-card'>"
"<h3>Break-glass "
f"{_bool_badge(bg.available, true_label='available', false_label='unavailable')}"
"</h3>"
f"<p class='muted'>{_esc_text(bg.reason)}</p>"
f"<p class='meta'>Tracked by issue #{_esc(bg.issue)}.</p>"
"</section>"
)
def render_restart_console_page(snapshot: RestartConsoleSnapshot) -> str:
"""Render the whole read-only restart console body."""
return (
"<h2>Restart status and impact</h2>"
f"<p class='meta'>Generated <code>{_esc(snapshot.generated_at)}</code> · "
f"viewer role <code>{_esc(snapshot.viewer_role)}</code> · "
f"authenticated {_bool_badge(snapshot.viewer_authenticated)} · "
f"read-only {_bool_badge(snapshot.read_only)}</p>"
+ _notes_block(snapshot)
+ _impact_section(snapshot)
+ _drain_section(snapshot)
+ _reconcile_section(snapshot)
+ _class_matrix_section(snapshot)
+ _approval_section(snapshot)
+ _break_glass_section(snapshot)
)
+47 -20
View File
@@ -270,26 +270,53 @@ def _probe_error_card(snapshot: SystemHealthSnapshot) -> str:
def _recovery_card() -> str:
"""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>"
)
"""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 &rarr; Preview &rarr; Confirm &rarr; 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>"
)
def render_system_health_page(snapshot: SystemHealthSnapshot) -> str:
+10
View File
@@ -201,6 +201,16 @@ 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.