Compare commits

..
Author SHA1 Message Date
sysadmin f0c6255d7d Merge remote-tracking branch 'prgs/master' into feat/issue-659-maintenance-drain-mode 2026-07-25 18:40:57 -04:00
sysadminandClaude Opus 4.8 77d808e7d4 merge(master): resolve PR #907 allocator drain vs side_effect_free
Keep #659 maintenance-drain assignment stop and #643 side_effect_free
apply guard; session registration remains gated by side_effect_free.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 18:12:03 -04:00
sysadminandClaude Opus 4.8 e91b94db56 feat(mcp): implement graceful maintenance-drain mode (Closes #659)
Add durable per-repo drain state, allocator assignment stop, mutation
deferral with a safety allowlist, observable status, and capability-gated
enter/exit tools. Drain proof/restart gate remain #661.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 17:02:12 -04:00
17 changed files with 1056 additions and 1200 deletions
+41
View File
@@ -26,6 +26,7 @@ from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Mapping, Sequence
import maintenance_drain
from control_plane_db import (
ControlPlaneDB,
ControlPlaneError,
@@ -936,6 +937,46 @@ def allocate_next_work(
"allocation_mode": (allocation_mode or "").strip() or None,
}
# #659 AC2: while maintenance drain is active, no new work is assigned —
# for dry-run and apply alike, so a preview can never be read as evidence
# that work was assignable during the drain. Checked before session
# registration so a drained allocator leaves no new state behind.
try:
drain_record = db.read_maintenance_drain(remote=remote, org=org, repo=repo)
except Exception as exc: # noqa: BLE001 — unreadable drain state fails closed
return {
"success": False,
"outcome": OUTCOME_NO_SAFE,
"reasons": [
f"maintenance-drain state lookup failed: {exc} (fail closed, #659)"
],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
}
drain_decision = maintenance_drain.classify_assignment(drain_record)
if not drain_decision["assignment_allowed"]:
return {
"success": True,
"outcome": OUTCOME_WAIT,
"apply": apply,
"role": role_norm,
"allocation_mode": mode,
"remote": remote,
"org": org,
"repo": repo,
"selected": None,
"reasons": list(drain_decision["reasons"]),
"reason_code": drain_decision["reason_code"],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
"maintenance_drain": maintenance_drain.status_payload(
drain_record, remote=remote, org=org, repo=repo
),
}
# 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:
+194 -27
View File
@@ -31,8 +31,9 @@ from typing import Any, Iterator, Sequence
import dependency_graph
import gitea_audit
import maintenance_drain
SCHEMA_VERSION = 5
SCHEMA_VERSION = 6
# Assignable work kinds only — raw monitoring incidents are never work items.
WORK_KINDS = frozenset({"issue", "pr"})
@@ -239,6 +240,31 @@ CREATE INDEX IF NOT EXISTS idx_session_checkpoints_session
CREATE INDEX IF NOT EXISTS idx_session_checkpoints_work
ON session_checkpoints(remote, org, repo, work_kind, work_number);
-- Graceful maintenance-drain state (#659). One current row per repository
-- scope — drain is a *state*, not a history, so entering and exiting update
-- the same row and every transition is audited to ``events``. Creating the
-- table is the v5->v6 migration: additive, idempotent, and it never touches
-- prior tables. ``state`` is CHECK-constrained so an unknown value can never
-- be written and later read as "not draining".
CREATE TABLE IF NOT EXISTS maintenance_drain (
drain_id TEXT PRIMARY KEY,
remote TEXT NOT NULL,
org TEXT NOT NULL,
repo TEXT NOT NULL,
state TEXT NOT NULL DEFAULT 'inactive'
CHECK (state IN ('inactive', 'draining')),
reason TEXT NOT NULL DEFAULT '',
requested_by TEXT NOT NULL DEFAULT '',
requested_by_profile TEXT NOT NULL DEFAULT '',
session_id TEXT NOT NULL DEFAULT '',
entered_at TEXT NOT NULL DEFAULT '',
exited_at TEXT NOT NULL DEFAULT '',
drain_schema_version INTEGER NOT NULL DEFAULT 6,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
UNIQUE (remote, org, repo)
);
-- Model usage, token cost, latency, and performance events (#651)
CREATE TABLE IF NOT EXISTS usage_events (
usage_id INTEGER PRIMARY KEY AUTOINCREMENT,
@@ -1568,32 +1594,6 @@ class ControlPlaneDB:
).fetchone()
return dict(row) if row else None
def list_incident_links(
self,
*,
provider: str | None = None,
gitea_org: str | None = None,
gitea_repo: str | None = None,
limit: int = 100,
) -> list[dict[str, Any]]:
"""List stored incident_links rows, optionally filtered by provider/repo (#612 / #649)."""
query = "SELECT * FROM incident_links WHERE 1=1"
params: list[Any] = []
if provider:
query += " AND provider = ?"
params.append(provider.strip().lower())
if gitea_org:
query += " AND gitea_org = ?"
params.append(_norm_scope(gitea_org))
if gitea_repo:
query += " AND gitea_repo = ?"
params.append(_norm_scope(gitea_repo))
query += " ORDER BY link_id DESC LIMIT ?"
params.append(max(1, limit))
with self._tx(immediate=False) as conn:
rows = conn.execute(query, params).fetchall()
return [dict(r) for r in rows]
# ── lease lifecycle (#601) ────────────────────────────────────────────
@@ -3051,3 +3051,170 @@ class ControlPlaneDB:
"live_lease_id": None if live_lease_id is None else str(live_lease_id),
"reconcile_action": "reconcile_required" if stale else "safe_to_resume",
}
# ── Maintenance drain (#659) ─────────────────────────────────────────────
@staticmethod
def _maintenance_drain_row(row: sqlite3.Row | None) -> dict[str, Any] | None:
"""Convert a ``maintenance_drain`` row to a plain record."""
if row is None:
return None
return {key: row[key] for key in row.keys()}
def read_maintenance_drain(
self, *, remote: str, org: str, repo: str
) -> dict[str, Any] | None:
"""Return the current drain record for a scope, or None if never set.
None and a stored ``inactive`` row mean the same thing to callers —
``maintenance_drain.is_draining`` treats both as not draining — so the
read never has to invent a record to answer the gate.
"""
with self._tx(immediate=False) as conn:
row = conn.execute(
"""
SELECT * FROM maintenance_drain
WHERE remote = ? AND org = ? AND repo = ?
""",
(str(remote or ""), str(org or ""), str(repo or "")),
).fetchone()
return self._maintenance_drain_row(row)
def set_maintenance_drain(
self,
*,
remote: str,
org: str,
repo: str,
state: str,
reason: str = "",
requested_by: str = "",
requested_by_profile: str = "",
session_id: str = "",
) -> dict[str, Any]:
"""Enter or exit maintenance drain for one repository scope (AC1).
The state transition is audited to ``events`` — entering and exiting
are exactly the moments an operator has to be able to reconstruct
later. Re-entering an already-draining scope is idempotent: it refreshes
the reason/owner metadata, keeps the original ``entered_at``, and
records no duplicate transition event.
Capability authorization happens above this layer (the drain tasks
carry a non-``gitea.*`` permission in the task capability map); the DB
records who asked and why, and never grants the right itself.
"""
state_norm = maintenance_drain.normalize_state(state)
raw = {
"reason": str(reason or ""),
"requested_by": str(requested_by or ""),
"requested_by_profile": str(requested_by_profile or ""),
"session_id": str(session_id or ""),
}
clean = gitea_audit.redact(raw)
remote_s, org_s, repo_s = str(remote or ""), str(org or ""), str(repo or "")
now_s = _ts()
with self._tx() as conn:
existing = conn.execute(
"""
SELECT * FROM maintenance_drain
WHERE remote = ? AND org = ? AND repo = ?
""",
(remote_s, org_s, repo_s),
).fetchone()
prior_state = (
maintenance_drain.normalize_state(existing["state"])
if existing is not None
else maintenance_drain.STATE_INACTIVE
)
transitioned = prior_state != state_norm
prior_entered = (
str(existing["entered_at"] or "") if existing is not None else ""
)
prior_exited = (
str(existing["exited_at"] or "") if existing is not None else ""
)
if state_norm == maintenance_drain.STATE_DRAINING:
# A re-entry keeps the original entry time (the drain never
# stopped); a fresh entry stamps now and clears the old exit.
entered_at = prior_entered if (not transitioned and prior_entered) else now_s
exited_at = ""
else:
entered_at = prior_entered
exited_at = now_s if (transitioned or not prior_exited) else prior_exited
if existing is None:
drain_id = uuid.uuid4().hex
conn.execute(
"""
INSERT INTO maintenance_drain(
drain_id, remote, org, repo, state, reason,
requested_by, requested_by_profile, session_id,
entered_at, exited_at, drain_schema_version,
created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
drain_id, remote_s, org_s, repo_s, state_norm,
clean["reason"], clean["requested_by"],
clean["requested_by_profile"], clean["session_id"],
entered_at, exited_at,
maintenance_drain.DRAIN_SCHEMA_VERSION, now_s, now_s,
),
)
else:
drain_id = str(existing["drain_id"])
conn.execute(
"""
UPDATE maintenance_drain
SET state = ?, reason = ?, requested_by = ?,
requested_by_profile = ?, session_id = ?,
entered_at = ?, exited_at = ?,
drain_schema_version = ?, updated_at = ?
WHERE drain_id = ?
""",
(
state_norm, clean["reason"], clean["requested_by"],
clean["requested_by_profile"], clean["session_id"],
entered_at, exited_at,
maintenance_drain.DRAIN_SCHEMA_VERSION, now_s, drain_id,
),
)
if transitioned:
event_type = (
"maintenance_drain_enter"
if state_norm == maintenance_drain.STATE_DRAINING
else "maintenance_drain_exit"
)
conn.execute(
"""
INSERT INTO events(work_item_id, event_type, message, created_at)
VALUES (NULL, ?, ?, ?)
""",
(
event_type,
f"drain {drain_id} scope {remote_s}/{org_s}/{repo_s} "
f"{prior_state} -> {state_norm} by "
f"{clean['requested_by'] or '(unknown)'} "
f"({clean['requested_by_profile'] or 'no profile'}); "
f"reason: {clean['reason'] or '(none)'}",
now_s,
),
)
row = conn.execute(
"SELECT * FROM maintenance_drain WHERE drain_id = ?", (drain_id,)
).fetchone()
record = self._maintenance_drain_row(row) or {}
return {
"record": record,
"drain_id": drain_id,
"state": state_norm,
"prior_state": prior_state,
"transitioned": transitioned,
}
+45
View File
@@ -0,0 +1,45 @@
# MCP maintenance-drain mode (#659)
Graceful **maintenance drain** stops new work assignment and defers non-allowlisted
mutations so sessions can finish critical handoffs and checkpoint before a
restart. It is **not** a restart authorization: the drain *proof* and apply gate
remain #661.
## State
Per repository scope (`remote`/`org`/`repo`) in the control-plane DB table
`maintenance_drain` (schema v6):
| State | Meaning |
|-------|---------|
| `inactive` | Normal operation (also: no row) |
| `draining` | Assignment stopped; non-allowlisted mutations deferred |
Enter/exit transitions are audited as `maintenance_drain_enter` /
`maintenance_drain_exit` events.
## Tools
| Tool | Permission | Effect |
|------|------------|--------|
| `gitea_maintenance_drain_status` | `gitea.read` | Observe drain (every session) |
| `gitea_enter_maintenance_drain` | `runtime.maintenance_drain` | Enter drain (capability-gated) |
| `gitea_exit_maintenance_drain` | `runtime.maintenance_drain` | Exit drain |
`runtime.maintenance_drain` is intentionally **not** a `gitea.*` op, so ordinary
author profiles cannot enter drain by accident.
## Enforcement
1. **Allocator** (`allocate_next_work`): while draining, returns `outcome=wait`
with `reason_code=maintenance_drain_assignment_stopped` for dry-run and apply.
2. **Mutation preflight** (`verify_preflight_purity`): non-allowlisted mutation
tasks raise `MaintenanceDrainError` with a typed next action.
3. **Allowlist** (safety only): heartbeats, lease release/abandon, session
checkpoints, enter/exit drain. Reads always work.
## Restart relationship
Drain mode prepares the blast radius. Restart apply still requires a clean
`DrainProof` (#661) or authorized break-glass. Status payloads never claim
restart permission.
@@ -1,35 +0,0 @@
# Web Console: Sentry/GlitchTip Observability & Incident Bridge Console (#649)
This document describes the Phase 4 observability console surface integrated into the MCP Control Plane Web Console (`webui/`), backed by the #612 incident bridge and the #613 control-plane DB substrate.
## Architectural Authority Model (ADR Alignment)
Per the Web Console Architecture ADR (`docs/architecture/webui-control-plane-console-architecture-adr.md`):
| Layer | Responsibility | Authority |
|---|---|---|
| **Gitea** | Durable work record | Issues, PRs, comments, reviews, labels, merges |
| **Control-plane DB** | Live coordination & linkage | `incident_links` table, session leases, allocations |
| **Sentry / GlitchTip** | Observability input | Unresolved incidents, error events, stack traces |
| **Incident Bridge (#612)** | Reconciliation engine | Reconciles provider observations into Gitea issues |
| **Web Console (`webui/`)** | Read-only projection & gated actions | Projects connection health & correlation links; gates writes |
> **Key Rule:** Raw monitoring incidents are **never** assignable control-plane `work_items`. They remain observation input only.
## Redaction Boundary Invariants
1. **No secrets in returns or rendering:** Auth tokens (`SENTRY_AUTH_TOKEN`, `GLITCHTIP_AUTH_TOKEN`), DSNs, `Authorization` headers, and sensitive local file paths are passed through `webui.console_redaction` before leaving the server.
2. **Safe projection:** Connection objects report `credentials_present: true/false` rather than exposing raw keys or headers.
## Console Endpoints
- **HTML Surface:** `GET /observability` — Renders provider connection cards, error correlation tables, and gated reconcile controls.
- **Versioned API:** `GET /api/v1/observability` — Returns structured JSON snapshot with `schema_version`, `providers`, `links`, and `metrics`.
- **Legacy Compatibility Alias:** `GET /api/observability` — Read-only compatibility alias for Phase 4.
## Gated Actions
- `observability_reconcile_incident` (`gitea_observability_reconcile_incident`): Triggers or previews dry-run issue reconciliation for a provider incident.
- `observability_link_issue` (`gitea_observability_link_issue`): Links a provider incident to an existing Gitea tracking issue.
Both actions require `operator` role and gate through `task_capability_map`. Execution fails closed in read-only MVP mode.
-2
View File
@@ -98,8 +98,6 @@ already define, and a regression test asserts each mapping matches.
| `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 |
| `observability_reconcile_incident` | operator | gated_write | `gitea.read` | Yes | No | No | 4 |
| `observability_link_issue` | operator | gated_write | `gitea.read` | Yes | No | No | 4 |
**Dual control** means the acting principal may not be the sole authority: a
second distinct principal must confirm. **Break-glass** means the action is
+240
View File
@@ -1472,6 +1472,9 @@ def verify_preflight_purity(
# contaminated by manual MCP daemon process killing (reconciler-exempt).
_enforce_runtime_recovery_contamination_gate(task, remote)
# #659 AC3: defer non-allowlisted mutations while maintenance drain is active.
_enforce_maintenance_drain_gate(task, remote=remote, org=org, repo=repo)
ctx = _resolve_namespace_mutation_context(worktree_path)
workspace = ctx["workspace_path"]
canonical_root = ctx["canonical_repo_root"]
@@ -2004,6 +2007,55 @@ def _enforce_runtime_recovery_contamination_gate(
)
def _enforce_maintenance_drain_gate(
task: str | None,
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> None:
"""#659 AC3: defer non-allowlisted mutations while drain is active.
The single mutation chokepoint already used by every gated task, so drain
coverage cannot drift per-tool. Allowlisted safety operations (heartbeat,
release/abandon, checkpoint, drain exit) pass through so an in-flight
session can still finish and hand off; everything else is deferred with a
typed blocker. Unreadable drain state fails closed a drain that cannot be
read is not evidence that no drain is running.
"""
if _preflight_in_test_mode() and not os.environ.get(
"GITEA_TEST_FORCE_MAINTENANCE_DRAIN"
):
return
if maintenance_drain.is_allowlisted_task(task):
return
try:
_h, o, r = _resolve(remote, None, org, repo)
except Exception: # noqa: BLE001 — scope resolution is best-effort here
o, r = (org or ""), (repo or "")
db, errs = _control_plane_db_or_error()
if db is None:
raise RuntimeError(
"maintenance-drain state could not be read: "
f"{'; '.join(errs) or 'control-plane DB unavailable'} (fail closed, #659)"
)
try:
record = db.read_maintenance_drain(remote=remote or "", org=o, repo=r)
except Exception as exc: # noqa: BLE001
raise RuntimeError(
f"maintenance-drain state could not be read: {_redact(str(exc))} "
"(fail closed, #659)"
) from exc
decision = maintenance_drain.classify_mutation(task, record)
if not decision["allowed"]:
raise maintenance_drain.MaintenanceDrainError(
maintenance_drain.format_drain_block_error(decision),
decision=decision,
)
def _enforce_stable_branch_contamination_gate(
task: str | None,
remote: str | None = None,
@@ -2064,6 +2116,7 @@ import allocator_service # noqa: E402
import allocator_dependencies # noqa: E402
import dependency_graph # noqa: E402 # #784 durable dependency edges
import control_plane_db # noqa: E402
import maintenance_drain # noqa: E402 # #659 graceful maintenance-drain mode
import lease_lifecycle # noqa: E402
import lease_policy # noqa: E402
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
@@ -22559,6 +22612,193 @@ def gitea_workflow_dashboard(
return payload
@mcp.tool()
def gitea_maintenance_drain_status(
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Read-only: current maintenance-drain state for a repository scope (#659 AC4).
Every session must be able to observe drain so it can stop creating new work
and finish only allowlisted safety operations. Never mutates; never restarts.
"""
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"read_only": True,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
try:
_h, o, r = _resolve(remote, host, org, repo)
except ValueError as exc:
return {"success": False, "read_only": True, "reasons": [str(exc)]}
db, errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"read_only": True,
"reasons": errs or ["control-plane DB unavailable"],
"maintenance_drain": maintenance_drain.status_payload(
None, remote=remote, org=o, repo=r
),
}
try:
record = db.read_maintenance_drain(remote=remote, org=o, repo=r)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"read_only": True,
"reasons": [f"drain state unreadable: {_redact(str(exc))}"],
}
payload = maintenance_drain.status_payload(
record, remote=remote, org=o, repo=r
)
return {"success": True, "read_only": True, "maintenance_drain": payload}
@mcp.tool()
def gitea_enter_maintenance_drain(
reason: str = "",
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
session_id: str | None = None,
) -> dict:
"""Enter graceful maintenance-drain mode for a repository scope (#659 AC1).
Stops new assignment and defers non-allowlisted mutations until exit. Requires
``runtime.maintenance_drain`` (controller/lifecycle capability not granted
by ordinary Gitea author profiles). Audited in the control-plane event log.
"""
cap_block = _profile_operation_gate("runtime.maintenance_drain")
if cap_block:
return {
"success": False,
"performed": False,
"reasons": cap_block,
"permission_report": _permission_block_report(
"runtime.maintenance_drain"
),
}
try:
_h, o, r = _resolve(remote, host, org, repo)
except ValueError as exc:
return {"success": False, "performed": False, "reasons": [str(exc)]}
db, errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"performed": False,
"reasons": errs or ["control-plane DB unavailable"],
}
profile = get_profile() or {}
try:
result = db.set_maintenance_drain(
remote=remote,
org=o,
repo=r,
state=maintenance_drain.STATE_DRAINING,
reason=reason or "operator-entered maintenance drain",
requested_by=str(
(profile.get("identity") or {}).get("username")
or profile.get("expected_username")
or ""
),
requested_by_profile=str(profile.get("profile_name") or ""),
session_id=str(session_id or ""),
)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"performed": False,
"reasons": [f"enter drain failed: {_redact(str(exc))}"],
}
record = result.get("record") or {}
return {
"success": True,
"performed": True,
"transitioned": bool(result.get("transitioned")),
"state": result.get("state"),
"prior_state": result.get("prior_state"),
"drain_id": result.get("drain_id"),
"maintenance_drain": maintenance_drain.status_payload(
record, remote=remote, org=o, repo=r
),
}
@mcp.tool()
def gitea_exit_maintenance_drain(
reason: str = "",
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
session_id: str | None = None,
) -> dict:
"""Exit graceful maintenance-drain mode (#659 AC1). Restores assignment and mutations."""
cap_block = _profile_operation_gate("runtime.maintenance_drain")
if cap_block:
return {
"success": False,
"performed": False,
"reasons": cap_block,
"permission_report": _permission_block_report(
"runtime.maintenance_drain"
),
}
try:
_h, o, r = _resolve(remote, host, org, repo)
except ValueError as exc:
return {"success": False, "performed": False, "reasons": [str(exc)]}
db, errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"performed": False,
"reasons": errs or ["control-plane DB unavailable"],
}
profile = get_profile() or {}
try:
result = db.set_maintenance_drain(
remote=remote,
org=o,
repo=r,
state=maintenance_drain.STATE_INACTIVE,
reason=reason or "operator-exited maintenance drain",
requested_by=str(
(profile.get("identity") or {}).get("username")
or profile.get("expected_username")
or ""
),
requested_by_profile=str(profile.get("profile_name") or ""),
session_id=str(session_id or ""),
)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"performed": False,
"reasons": [f"exit drain failed: {_redact(str(exc))}"],
}
record = result.get("record") or {}
return {
"success": True,
"performed": True,
"transitioned": bool(result.get("transitioned")),
"state": result.get("state"),
"prior_state": result.get("prior_state"),
"drain_id": result.get("drain_id"),
"maintenance_drain": maintenance_drain.status_payload(
record, remote=remote, org=o, repo=r
),
}
@mcp.tool()
def gitea_request_mcp_restart(
remote: str = "dadeschools",
+281
View File
@@ -0,0 +1,281 @@
"""Graceful MCP maintenance-drain mode (#659).
Drain is the visible, capability-gated state that lets an operator stop new
work and quiesce mutations *before* a restart, instead of cutting sessions off
mid-mutation. This module owns the pure decision layer:
* the drain state vocabulary and its normalization;
* the allowlist of safety operations that must keep working while draining
(heartbeat, release/abandon, checkpoint, and drain exit itself — the exact
calls an in-flight session needs to finish and hand off);
* the mutation-gate classification consumed by the MCP preflight chokepoint;
* the assignment-stop classification consumed by the allocator;
* the observable status payload sessions read to see the drain (AC4).
Durable state lives in the control-plane DB (``maintenance_drain`` table);
enforcement lives at the existing chokepoints. Nothing here performs I/O, so
both callers can share one decision without importing each other.
Scope note: the machine-verifiable *drain proof* and the restart gate that
consumes it are #661's scope, not this module's. Drain here stops assignment
and mutation and makes the state observable; it never authorizes a restart.
"""
from __future__ import annotations
from typing import Any, Mapping
# ── State vocabulary ──────────────────────────────────────────────────────────
STATE_INACTIVE = "inactive"
STATE_DRAINING = "draining"
DRAIN_STATES = frozenset({STATE_INACTIVE, STATE_DRAINING})
# Typed blocker code surfaced to clients (never a bare string at call sites).
BLOCKER_DRAIN_ACTIVE = "maintenance_drain_active"
# Reason code for the allocator's assignment stop.
REASON_ASSIGNMENT_STOPPED = "maintenance_drain_assignment_stopped"
DRAIN_SCHEMA_VERSION = 6
class MaintenanceDrainError(RuntimeError):
"""Raised when a mutation is refused because drain is active (fail closed)."""
def __init__(self, message: str, *, decision: Mapping[str, Any] | None = None):
super().__init__(message)
self.decision = dict(decision or {})
self.reason_code = BLOCKER_DRAIN_ACTIVE
# ── Safety allowlist ──────────────────────────────────────────────────────────
# Mutations that stay permitted while draining. Every entry is a *quiesce*
# operation: it either proves an in-flight task is still alive, hands its claim
# back, records the durable state a restart needs, or ends the drain. Nothing
# that creates new work, new branches, new PRs, or new review/merge verdicts is
# on this list — that is the whole point of the drain.
ALLOWLISTED_DRAIN_TASKS: frozenset[str] = frozenset(
{
# Liveness of work already in flight.
"heartbeat_issue_lock",
"heartbeat_reviewer_pr_lease",
"post_heartbeat",
# Handing claims back so nothing is stranded across the restart.
"release_workflow_lease",
"release_reviewer_pr_lease",
"release_merger_pr_lease",
"abandon_workflow_lease",
# Durable recovery state (#660) must be writable *during* drain.
"write_session_checkpoint",
"checkpoint_session",
# The drain controls themselves — exit must never be self-blocked.
"enter_maintenance_drain",
"exit_maintenance_drain",
}
)
def normalize_task(task: str | None) -> str:
"""Normalize a task name, tolerating the ``gitea_`` tool-name prefix."""
name = str(task or "").strip()
if name.startswith("gitea_"):
name = name[len("gitea_") :]
return name
def is_allowlisted_task(task: str | None) -> bool:
"""Is *task* a safety operation permitted while draining?"""
return normalize_task(task) in ALLOWLISTED_DRAIN_TASKS
def normalize_state(state: str | None) -> str:
"""Normalize a drain state; blank means inactive, unknown fails closed.
Blank normalizes to ``inactive`` (no drain record = not draining), but an
unrecognized non-blank value raises: silently treating ``"drainig"`` as
inactive would disable the gate.
"""
value = str(state or "").strip().lower()
if not value:
return STATE_INACTIVE
if value not in DRAIN_STATES:
raise MaintenanceDrainError(
f"unknown maintenance-drain state {value!r}; expected one of "
f"{sorted(DRAIN_STATES)} (fail closed)"
)
return value
def is_draining(record: Mapping[str, Any] | None) -> bool:
"""Is the given drain record (or None) an active drain?"""
if not record:
return False
return normalize_state(record.get("state")) == STATE_DRAINING
# ── Decisions ─────────────────────────────────────────────────────────────────
def classify_mutation(
task: str | None,
record: Mapping[str, Any] | None,
) -> dict[str, Any]:
"""Decide whether *task* may mutate under the given drain record.
Returns a decision dict with ``allowed``/``deferred`` and, when refused, a
typed ``reason_code`` plus the one exact next action the caller may take.
Deferred (not failed): the operation is legal again after drain exits, so
the caller is told to wait rather than to retry a different way.
"""
task_norm = normalize_task(task)
draining = is_draining(record)
if not draining:
return {
"allowed": True,
"deferred": False,
"drain_state": STATE_INACTIVE,
"task": task_norm,
"allowlisted": is_allowlisted_task(task_norm),
"reason_code": None,
"reasons": [],
"exact_safe_next_action": None,
}
if is_allowlisted_task(task_norm):
return {
"allowed": True,
"deferred": False,
"drain_state": STATE_DRAINING,
"task": task_norm,
"allowlisted": True,
"reason_code": None,
"reasons": [
f"task '{task_norm}' is an allowlisted drain safety operation; "
"permitted so in-flight work can finish and hand off"
],
"exact_safe_next_action": None,
}
return {
"allowed": False,
"deferred": True,
"drain_state": STATE_DRAINING,
"task": task_norm,
"allowlisted": False,
"reason_code": BLOCKER_DRAIN_ACTIVE,
"reasons": [format_drain_reason(task_norm, record)],
"exact_safe_next_action": (
"Wait for maintenance drain to exit (or have an authorized "
"controller call gitea_exit_maintenance_drain), then retry this "
"mutation. Reads and gitea_maintenance_drain_status stay available."
),
}
def classify_assignment(record: Mapping[str, Any] | None) -> dict[str, Any]:
"""Decide whether the allocator may assign new work (AC2)."""
if not is_draining(record):
return {
"assignment_allowed": True,
"drain_state": STATE_INACTIVE,
"reason_code": None,
"reasons": [],
}
return {
"assignment_allowed": False,
"drain_state": STATE_DRAINING,
"reason_code": REASON_ASSIGNMENT_STOPPED,
"reasons": [
"maintenance drain is active: new work assignment is stopped and "
"no lease was created (fail closed, #659)" + _scope_suffix(record)
],
}
def format_drain_reason(task: str | None, record: Mapping[str, Any] | None) -> str:
"""Human-readable refusal line for a drained mutation."""
task_norm = normalize_task(task) or "(unnamed task)"
return (
f"maintenance drain is active: mutation '{task_norm}' is deferred; only "
"allowlisted drain safety operations "
f"({', '.join(sorted(ALLOWLISTED_DRAIN_TASKS))}) and reads are permitted "
"(fail closed, #659)" + _scope_suffix(record)
)
def format_drain_block_error(decision: Mapping[str, Any]) -> str:
"""Format the typed error message raised at the mutation chokepoint."""
reasons = list(decision.get("reasons") or [])
head = reasons[0] if reasons else "maintenance drain is active (fail closed)"
action = decision.get("exact_safe_next_action")
return f"{head}. Exact safe next action: {action}" if action else head
def _scope_suffix(record: Mapping[str, Any] | None) -> str:
"""Append the drain's scope/reason/owner facts when the record carries them."""
if not record:
return ""
bits: list[str] = []
scope = "/".join(
str(record.get(key) or "") for key in ("remote", "org", "repo")
).strip("/")
if scope:
bits.append(f"scope {scope}")
if record.get("reason"):
bits.append(f"reason: {record['reason']}")
if record.get("requested_by"):
bits.append(f"entered by {record['requested_by']}")
if record.get("entered_at"):
bits.append(f"at {record['entered_at']}")
return f" ({'; '.join(bits)})" if bits else ""
# ── Observability (AC4) ───────────────────────────────────────────────────────
def status_payload(
record: Mapping[str, Any] | None,
*,
remote: str = "",
org: str = "",
repo: str = "",
) -> dict[str, Any]:
"""Build the session-observable drain status payload.
Always answers, including when no drain record exists: an absent record is
a definitive "not draining", not an unknown.
"""
draining = is_draining(record)
rec: Mapping[str, Any] = record or {}
return {
"drain_state": STATE_DRAINING if draining else STATE_INACTIVE,
"draining": draining,
"remote": str(rec.get("remote") or "") or remote,
"org": str(rec.get("org") or "") or org,
"repo": str(rec.get("repo") or "") or repo,
"reason": str(rec.get("reason") or ""),
"requested_by": str(rec.get("requested_by") or ""),
"requested_by_profile": str(rec.get("requested_by_profile") or ""),
"session_id": str(rec.get("session_id") or ""),
"entered_at": str(rec.get("entered_at") or ""),
"exited_at": str(rec.get("exited_at") or ""),
"assignment_stopped": draining,
"mutations_deferred": draining,
"allowlisted_tasks": sorted(ALLOWLISTED_DRAIN_TASKS),
"reads_permitted": True,
"record_present": bool(record),
"schema_version": DRAIN_SCHEMA_VERSION,
"drain_proof_scope": (
"drain proof and the restart gate that consumes it are #661 scope; "
"this status never authorizes a restart"
),
"safe_next_action": (
"Wait for drain to exit before retrying deferred mutations; "
"allowlisted safety operations and reads remain available."
if draining
else "None; maintenance drain is not active."
),
}
+18
View File
@@ -414,6 +414,24 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"role": "controller",
},
# #659 maintenance drain. Same reasoning as the lifecycle controls above:
# entering/exiting drain quiesces a whole namespace, so it carries a
# non-``gitea.*`` permission that no configured Gitea profile satisfies by
# accident (AC1 — capability-gated and audited). Reading drain state is
# ordinary read authority: every session must be able to see the drain (AC4).
"enter_maintenance_drain": {
"permission": "runtime.maintenance_drain",
"role": "controller",
},
"exit_maintenance_drain": {
"permission": "runtime.maintenance_drain",
"role": "controller",
},
"maintenance_drain_status": {
"permission": "gitea.read",
"role": "author",
},
# #601 first-class lease lifecycle — inspect/list need read; mutations gate on
# ownership in the control-plane DB (not a separate Gitea write permission).
"list_workflow_leases": {
+237
View File
@@ -0,0 +1,237 @@
"""Tests for graceful MCP maintenance-drain mode (#659).
Acceptance coverage:
1. Enter/exit is durable and audited (DB substrate).
2. New work assignment stops during drain (allocator WAIT).
3. Mutations deferred except allowlisted safety ops.
4. Sessions can observe drain state.
5. Fail-closed on unreadable drain state.
"""
from __future__ import annotations
import sys
import tempfile
import unittest
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
import maintenance_drain
from control_plane_db import ControlPlaneDB
from allocator_service import WorkCandidate, allocate_next_work, OUTCOME_WAIT
class TestDrainDecisions(unittest.TestCase):
def test_inactive_allows_mutations_and_assignment(self):
decision = maintenance_drain.classify_mutation("create_pr", None)
self.assertTrue(decision["allowed"])
self.assertFalse(decision["deferred"])
assign = maintenance_drain.classify_assignment(None)
self.assertTrue(assign["assignment_allowed"])
def test_draining_defers_non_allowlisted_mutation(self):
record = {"state": "draining", "remote": "prgs", "org": "o", "repo": "r"}
decision = maintenance_drain.classify_mutation("create_pr", record)
self.assertFalse(decision["allowed"])
self.assertTrue(decision["deferred"])
self.assertEqual(decision["reason_code"], maintenance_drain.BLOCKER_DRAIN_ACTIVE)
self.assertIn("create_pr", decision["reasons"][0])
def test_allowlisted_safety_ops_pass_during_drain(self):
record = {"state": "draining"}
for task in (
"heartbeat_issue_lock",
"gitea_release_reviewer_pr_lease",
"write_session_checkpoint",
"exit_maintenance_drain",
):
with self.subTest(task=task):
decision = maintenance_drain.classify_mutation(task, record)
self.assertTrue(decision["allowed"], decision)
def test_assignment_stopped_during_drain(self):
record = {"state": "draining", "reason": "upgrade"}
decision = maintenance_drain.classify_assignment(record)
self.assertFalse(decision["assignment_allowed"])
self.assertEqual(
decision["reason_code"], maintenance_drain.REASON_ASSIGNMENT_STOPPED
)
def test_unknown_state_fails_closed(self):
with self.assertRaises(maintenance_drain.MaintenanceDrainError):
maintenance_drain.normalize_state("drainig")
def test_status_payload_always_answers(self):
inactive = maintenance_drain.status_payload(None, remote="prgs", org="o", repo="r")
self.assertFalse(inactive["draining"])
self.assertTrue(inactive["reads_permitted"])
active = maintenance_drain.status_payload(
{"state": "draining", "reason": "reboot", "requested_by": "ops"},
remote="prgs",
org="o",
repo="r",
)
self.assertTrue(active["draining"])
self.assertTrue(active["assignment_stopped"])
self.assertTrue(active["mutations_deferred"])
self.assertIn("heartbeat_issue_lock", active["allowlisted_tasks"])
class TestDrainDB(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.db = ControlPlaneDB(db_path=str(Path(self._tmp.name) / "cp.sqlite3"))
def test_enter_exit_idempotent_and_audited(self):
first = self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="draining",
reason="planned restart",
requested_by="sysadmin",
requested_by_profile="prgs-controller",
session_id="s1",
)
self.assertTrue(first["transitioned"])
self.assertEqual(first["state"], "draining")
self.assertTrue(maintenance_drain.is_draining(first["record"]))
again = self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="draining",
reason="still draining",
requested_by="sysadmin",
requested_by_profile="prgs-controller",
session_id="s1",
)
self.assertFalse(again["transitioned"])
self.assertEqual(again["record"]["entered_at"], first["record"]["entered_at"])
exited = self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="inactive",
reason="done",
requested_by="sysadmin",
requested_by_profile="prgs-controller",
session_id="s1",
)
self.assertTrue(exited["transitioned"])
self.assertFalse(maintenance_drain.is_draining(exited["record"]))
self.assertTrue(exited["record"]["exited_at"])
# Events recorded for transitions only (enter + exit).
with self.db._tx(immediate=False) as conn:
rows = conn.execute(
"SELECT event_type FROM events WHERE event_type LIKE 'maintenance_drain_%' "
"ORDER BY event_id"
).fetchall()
types = [r[0] for r in rows]
self.assertEqual(types, ["maintenance_drain_enter", "maintenance_drain_exit"])
def test_read_missing_is_none_not_error(self):
self.assertIsNone(
self.db.read_maintenance_drain(remote="prgs", org="o", repo="r")
)
class TestAllocatorStopsDuringDrain(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.db = ControlPlaneDB(db_path=str(Path(self._tmp.name) / "cp.sqlite3"))
def test_allocate_returns_wait_while_draining(self):
self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="draining",
reason="test",
requested_by="tester",
)
candidates = [
WorkCandidate(
kind="issue",
number=659,
title="drain",
labels=("status:ready",),
priority=20,
)
]
result = allocate_next_work(
self.db,
role="author",
session_id="test-session",
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
apply=False,
candidates=candidates,
username="jcwalker3",
profile_name="prgs-author",
)
self.assertEqual(result["outcome"], OUTCOME_WAIT)
self.assertIsNone(result.get("selected"))
self.assertEqual(
result.get("reason_code"),
maintenance_drain.REASON_ASSIGNMENT_STOPPED,
)
self.assertTrue(result["maintenance_drain"]["draining"])
def test_allocate_works_when_inactive(self):
candidates = [
WorkCandidate(
kind="issue",
number=659,
title="drain",
labels=("status:ready",),
priority=20,
)
]
result = allocate_next_work(
self.db,
role="author",
session_id="test-session-2",
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
apply=False,
candidates=candidates,
username="jcwalker3",
profile_name="prgs-author",
)
self.assertNotEqual(
result.get("reason_code"),
maintenance_drain.REASON_ASSIGNMENT_STOPPED,
)
class TestCapabilityMap(unittest.TestCase):
def test_drain_tasks_mapped(self):
import task_capability_map as tcm
self.assertEqual(
tcm.required_permission("enter_maintenance_drain"),
"runtime.maintenance_drain",
)
self.assertEqual(
tcm.required_permission("exit_maintenance_drain"),
"runtime.maintenance_drain",
)
self.assertEqual(
tcm.required_permission("maintenance_drain_status"),
"gitea.read",
)
if __name__ == "__main__":
unittest.main()
-478
View File
@@ -1,478 +0,0 @@
"""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()
-190
View File
@@ -1,190 +0,0 @@
"""Tests for Sentry/GlitchTip observability console (#649, Phase 4)."""
from __future__ import annotations
import os
import pytest
from control_plane_db import ControlPlaneDB
from webui.app import create_app
from webui.console_authz import authorize, resolve_principal
from webui.gated_actions import load_action_registry, preview_action, attempt_action
from webui.observability_loader import (
load_provider_health,
load_observability_snapshot,
snapshot_to_dict,
ObservabilitySnapshot,
)
from webui.observability_views import render_observability_page
from tests.webui_testclient import TestClient
@pytest.fixture
def test_db(tmp_path):
db_path = str(tmp_path / "test_control_plane.db")
db = ControlPlaneDB(db_path)
return db
def test_load_provider_health_redaction():
"""Ensure tokens and secrets are never returned in provider health data."""
env = {
"SENTRY_BASE_URL": "https://sentry.prgs.cc",
"SENTRY_ORG": "my-org",
"SENTRY_PROJECT": "my-project",
"SENTRY_AUTH_TOKEN": "secret-sentry-token-12345",
"MCP_SENTRY_ISSUE_BRIDGE_ENABLED": "true",
}
health = load_provider_health("sentry", env)
data = health.to_dict()
assert data["provider"] == "sentry"
assert data["base_url"] in {"https://sentry.prgs.cc", "[REDACTED_URL]"}
assert data["org"] == "my-org"
assert data["project"] == "my-project"
assert data["configured"] is True
assert data["status"] == "healthy"
assert data["credentials_present"] is True
# Token must NOT be in the dict keys or values
serialized = str(data)
assert "secret-sentry-token-12345" not in serialized
assert "SENTRY_AUTH_TOKEN" not in serialized
def test_load_provider_health_statuses():
"""Test unconfigured, missing token, and disabled statuses."""
# Not configured
h1 = load_provider_health("sentry", {})
d1 = h1.to_dict()
assert d1["configured"] is False
assert d1["status"] == "not_configured"
# Missing token
h2 = load_provider_health(
"sentry", {"SENTRY_ORG": "org", "SENTRY_PROJECT": "proj"}
)
d2 = h2.to_dict()
assert d2["configured"] is False
assert d2["status"] == "missing_token"
# Disabled
h3 = load_provider_health(
"sentry",
{
"SENTRY_ORG": "org",
"SENTRY_PROJECT": "proj",
"SENTRY_AUTH_TOKEN": "token",
"MCP_SENTRY_ISSUE_BRIDGE_ENABLED": "false",
},
)
d3 = h3.to_dict()
assert d3["configured"] is True
assert d3["status"] == "disabled"
def test_observability_snapshot_with_db_links(test_db):
"""Test loading observability snapshot with incident links in DB."""
test_db.upsert_incident_link(
provider="sentry",
provider_issue_id="101",
gitea_org="Scaled-Tech-Consulting",
gitea_repo="Gitea-Tools",
gitea_issue_number=649,
provider_base_url="https://sentry.prgs.cc",
provider_org="Scaled-Tech-Consulting",
provider_project="Gitea-Tools",
provider_short_id="ST-101",
provider_permalink="https://sentry.prgs.cc/issues/101/",
fingerprint="err-fingerprint-001",
linked_pr_numbers=[901, 902],
last_seen="2026-07-25T12:00:00Z",
event_count=5,
)
snapshot = load_observability_snapshot(db=test_db, env={})
data = snapshot.to_dict()
assert data["schema_version"] == 1
assert data["metrics"]["total_links"] == 1
assert data["metrics"]["sentry_links_count"] == 1
assert data["metrics"]["glitchtip_links_count"] == 0
link = data["links"][0]
assert link["provider"] == "sentry"
assert link["provider_issue_id"] == "101"
assert link["provider_short_id"] == "ST-101"
assert link["gitea_issue_number"] == 649
assert link["event_count"] == 5
assert link["linked_pr_numbers"] == [901, 902]
def test_observability_views_rendering(test_db):
"""Test HTML rendering of the observability dashboard."""
snapshot = load_observability_snapshot(db=test_db, env={})
html_output = render_observability_page(snapshot)
assert "Observability &amp; Incident Bridge (#649)" in html_output or "Observability & Incident Bridge (#649)" in html_output or "Observability" in html_output
assert "ADR Authority Model:" in html_output
assert "Provider Connections" in html_output
assert "Correlated Incidents" in html_output
def test_webui_observability_routes():
"""Test Starlette HTTP routes for /observability and /api/v1/observability."""
client = TestClient(create_app())
# HTML page route
res_html = client.get("/observability")
assert res_html.status_code == 200
assert "text/html" in res_html.headers["content-type"]
assert "Observability" in res_html.text
# Versioned API route
res_api_v1 = client.get("/api/v1/observability")
assert res_api_v1.status_code == 200
assert "application/json" in res_api_v1.headers["content-type"]
data_v1 = res_api_v1.json()
assert "schema_version" in data_v1
assert "providers" in data_v1
assert "links" in data_v1
assert "metrics" in data_v1
# Compatibility alias route
res_api_alias = client.get("/api/observability")
assert res_api_alias.status_code == 200
assert res_api_alias.json() == data_v1
def test_observability_gated_actions():
"""Ensure observability actions are registered, gated, and fail closed in MVP mode."""
registry = load_action_registry()
action_reconcile = registry.get("observability_reconcile_incident")
assert action_reconcile is not None
assert action_reconcile.task_key == "observability_reconcile_incident"
assert action_reconcile.mcp_tool == "gitea_observability_reconcile_incident"
action_link = registry.get("observability_link_issue")
assert action_link is not None
assert action_link.task_key == "observability_link_issue"
# Preview returns mutation ledger
prev = preview_action("observability_reconcile_incident", provider="sentry", issue_id="101")
assert prev["action_id"] == "observability_reconcile_incident"
assert prev["enabled"] is False
# Execution fails closed in MVP mode
att = attempt_action("observability_reconcile_incident", provider="sentry", issue_id="101")
assert att["success"] is False
assert att["error"] == "action_disabled"
def test_observability_authz_rbac():
"""Test RBAC authorization for observability actions."""
principal = resolve_principal({})
# Check authorize decision
decision = authorize("observability_reconcile_incident", principal)
assert decision.action_id == "observability_reconcile_incident"
# Phase 4 action denies in Phase 1 runtime by default
assert decision.allowed is False
-21
View File
@@ -87,11 +87,6 @@ from webui.notifications import (
from webui.notification_views import render_notifications_page
from webui import request_service
from webui.request_views import render_requests_page
from webui.observability_loader import (
load_observability_snapshot,
snapshot_to_dict as observability_snapshot_to_dict,
)
from webui.observability_views import render_observability_page
_READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"})
_AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"})
@@ -915,19 +910,6 @@ async def api_notifications(request: Request) -> JSONResponse:
data = notifications_snapshot_to_dict(snap)
return JSONResponse(data)
async def observability_route(request: Request) -> HTMLResponse:
snap = load_observability_snapshot()
html_content = render_observability_page(snap)
return HTMLResponse(html_content)
async def api_observability(request: Request) -> JSONResponse:
snap = load_observability_snapshot()
data = observability_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.
@@ -1093,9 +1075,6 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/api/analytics", api_v1_analytics, methods=["GET"]),
Route("/api/v1/analytics", api_v1_analytics, methods=["GET"]),
Route("/api/v1/analytics/usage", api_v1_analytics_ingest, methods=["POST"]),
Route("/observability", observability_route, methods=["GET"]),
Route("/api/observability", api_observability, methods=["GET"]),
Route("/api/v1/observability", api_observability, methods=["GET"]),
Route("/audit", audit, methods=["GET", "POST"]),
Route("/api/audit", api_audit, methods=["GET", "POST"]),
Route("/worktrees", worktrees, methods=["GET"]),
-23
View File
@@ -317,29 +317,6 @@ _ACTION_SPECS: tuple[ConsoleAction, ...] = (
phase=2,
summary="Run reconciler cleanup for merged or superseded PR branches.",
),
# #649: Phase 4 observability & incident bridge actions.
ConsoleAction(
action_id="observability_reconcile_incident",
task_key="observability_reconcile_incident",
action_class=CLASS_WRITE,
minimum_role=OPERATOR,
requires_confirmation=True,
dual_control=False,
break_glass=False,
phase=4,
summary="Trigger/reconcile durable Gitea issue creation from a provider incident.",
),
ConsoleAction(
action_id="observability_link_issue",
task_key="observability_link_issue",
action_class=CLASS_WRITE,
minimum_role=OPERATOR,
requires_confirmation=True,
dual_control=False,
break_glass=False,
phase=4,
summary="Link a provider incident to an existing Gitea issue.",
),
# #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
-5
View File
@@ -185,11 +185,6 @@ def build_action_registry() -> ActionRegistry:
"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."),
# #649: Phase 4 observability & incident bridge actions.
("observability_reconcile_incident", "Reconcile incident", "observability_reconcile_incident",
"gitea_observability_reconcile_incident", "Trigger or dry-run durable issue reconciliation for a provider incident."),
("observability_link_issue", "Link incident issue", "observability_link_issue",
"gitea_observability_link_issue", "Link a provider incident to a Gitea tracking issue."),
)
actions = tuple(
GatedAction(
-1
View File
@@ -73,7 +73,6 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
)),
NavGroup("Insights", (
NavItem("/insights", "Insights", "stub"),
NavItem("/observability", "Observability"),
NavItem("/analytics", "Analytics"),
NavItem("/audit", "Audit"),
)),
-275
View File
@@ -1,275 +0,0 @@
"""Sentry/GlitchTip observability and incident correlation loader for the console (#649, Phase 4).
Operators need to inspect provider connection status (Sentry/GlitchTip), error
correlations, and durable Gitea issue linkage without treating raw incidents
as allocator work.
ADR authority model:
* Gitea owns work.
* Providers (Sentry/GlitchTip) observe incidents.
* Control-plane DB coordinates incident links.
* The #612 bridge reconciles observations into durable Gitea issues.
* The web console projects read-only state and gates mutations.
Redaction boundary:
* Provider auth tokens, DSNs, Authorization headers, and sensitive local file
paths are ALWAYS redacted before leaving this module.
"""
from __future__ import annotations
import os
from dataclasses import dataclass
from typing import Any
from control_plane_db import ControlPlaneDB
import sentry_incident_bridge
from webui import console_redaction
OBSERVABILITY_SCHEMA_VERSION = 1
@dataclass(frozen=True)
class ProviderHealth:
"""Connection and health status of an observability provider."""
provider: str
base_url: str
org: str
project: str
configured: bool
status: str
bridge_enabled: bool
lookback: str
min_events_for_issue: int
self_hosted: bool
environment: str | None = None
credentials_present: bool = False
def to_dict(self) -> dict[str, Any]:
data = {
"provider": self.provider,
"base_url": self.base_url,
"org": self.org,
"project": self.project,
"configured": self.configured,
"status": self.status,
"bridge_enabled": self.bridge_enabled,
"lookback": self.lookback,
"min_events_for_issue": self.min_events_for_issue,
"self_hosted": self.self_hosted,
"environment": self.environment,
"credentials_present": self.credentials_present,
}
return console_redaction.redact_payload(data)
@dataclass(frozen=True)
class CorrelatedIncidentLink:
"""One linked provider incident ↔ Gitea issue correlation record."""
link_id: int
provider: str
provider_base_url: str
provider_org: str
provider_project: str
provider_issue_id: str
provider_short_id: str | None
provider_permalink: str | None
fingerprint: str | None
gitea_org: str
gitea_repo: str
gitea_issue_number: int
linked_pr_numbers: list[int]
last_seen: str | None
event_count: int
created_at: str | None
updated_at: str | None
def to_dict(self) -> dict[str, Any]:
data = {
"link_id": self.link_id,
"provider": self.provider,
"provider_base_url": self.provider_base_url,
"provider_org": self.provider_org,
"provider_project": self.provider_project,
"provider_issue_id": self.provider_issue_id,
"provider_short_id": self.provider_short_id,
"provider_permalink": self.provider_permalink,
"fingerprint": self.fingerprint,
"gitea_org": self.gitea_org,
"gitea_repo": self.gitea_repo,
"gitea_issue_number": self.gitea_issue_number,
"linked_pr_numbers": self.linked_pr_numbers,
"last_seen": self.last_seen,
"event_count": self.event_count,
"created_at": self.created_at,
"updated_at": self.updated_at,
}
return console_redaction.redact_payload(data)
@dataclass(frozen=True)
class ObservabilitySnapshot:
"""Read-only snapshot of observability provider status and incident correlations."""
schema_version: int
providers: list[ProviderHealth]
links: list[CorrelatedIncidentLink]
total_links: int
sentry_links_count: int
glitchtip_links_count: int
bridge_active: bool
def to_dict(self) -> dict[str, Any]:
return {
"schema_version": self.schema_version,
"providers": [p.to_dict() for p in self.providers],
"links": [link.to_dict() for link in self.links],
"metrics": {
"total_links": self.total_links,
"sentry_links_count": self.sentry_links_count,
"glitchtip_links_count": self.glitchtip_links_count,
"bridge_active": self.bridge_active,
},
}
def load_provider_health(
provider_name: str = "sentry",
env: dict[str, str] | None = None,
) -> ProviderHealth:
"""Inspect configuration and connection health for an observability provider."""
source_env = dict(env if env is not None else os.environ)
if provider_name.lower() == "sentry":
config = sentry_incident_bridge.load_bridge_config(source_env)
token = sentry_incident_bridge.resolve_token(source_env)
has_token = bool(token)
configured = bool(config.org and config.project and has_token)
if not config.org or not config.project:
status = "not_configured"
elif not has_token:
status = "missing_token"
elif not config.bridge_enabled:
status = "disabled"
else:
status = "healthy"
return ProviderHealth(
provider="sentry",
base_url=config.base_url,
org=config.org or "unconfigured",
project=config.project or "unconfigured",
configured=configured,
status=status,
bridge_enabled=config.bridge_enabled,
lookback=config.lookback,
min_events_for_issue=config.min_events_for_issue,
self_hosted=not config.base_url.rstrip("/").endswith("sentry.io"),
environment=config.environment,
credentials_present=has_token,
)
# GlitchTip or fallback provider configuration
glitchtip_url = (source_env.get("GLITCHTIP_BASE_URL") or "https://glitchtip.prgs.cc").strip()
glitchtip_org = (source_env.get("GLITCHTIP_ORG") or "").strip()
glitchtip_proj = (source_env.get("GLITCHTIP_PROJECT") or "").strip()
glitchtip_token = (source_env.get("GLITCHTIP_AUTH_TOKEN") or "").strip()
has_token = bool(glitchtip_token)
configured = bool(glitchtip_org and glitchtip_proj and has_token)
status = "healthy" if configured else ("missing_token" if glitchtip_org and glitchtip_proj else "not_configured")
return ProviderHealth(
provider="glitchtip",
base_url=glitchtip_url,
org=glitchtip_org or "unconfigured",
project=glitchtip_proj or "unconfigured",
configured=configured,
status=status,
bridge_enabled=configured,
lookback="24h",
min_events_for_issue=2,
self_hosted=True,
environment=source_env.get("GLITCHTIP_ENVIRONMENT"),
credentials_present=has_token,
)
def _parse_pr_numbers(raw: Any) -> list[int]:
if isinstance(raw, list):
return [int(x) for x in raw if str(x).isdigit()]
if isinstance(raw, str) and raw.strip():
import json
try:
parsed = json.loads(raw)
if isinstance(parsed, list):
return [int(x) for x in parsed if str(x).isdigit()]
except Exception:
pass
return []
def load_observability_snapshot(
db: ControlPlaneDB | None = None,
env: dict[str, str] | None = None,
) -> ObservabilitySnapshot:
"""Build a read-only snapshot of observability connection health and incident links."""
sentry_health = load_provider_health("sentry", env)
glitchtip_health = load_provider_health("glitchtip", env)
providers = [sentry_health, glitchtip_health]
target_db = db or ControlPlaneDB()
raw_links = target_db.list_incident_links(limit=100)
links: list[CorrelatedIncidentLink] = []
sentry_cnt = 0
glitchtip_cnt = 0
for r in raw_links:
prov = (r.get("provider") or "sentry").lower()
if prov == "sentry":
sentry_cnt += 1
elif prov == "glitchtip":
glitchtip_cnt += 1
pr_nums = _parse_pr_numbers(r.get("linked_pr_numbers"))
links.append(
CorrelatedIncidentLink(
link_id=int(r.get("link_id", 0)),
provider=prov,
provider_base_url=r.get("provider_base_url") or "",
provider_org=r.get("provider_org") or "",
provider_project=r.get("provider_project") or "",
provider_issue_id=str(r.get("provider_issue_id") or ""),
provider_short_id=r.get("provider_short_id"),
provider_permalink=r.get("provider_permalink"),
fingerprint=r.get("fingerprint"),
gitea_org=r.get("gitea_org") or "Scaled-Tech-Consulting",
gitea_repo=r.get("gitea_repo") or "Gitea-Tools",
gitea_issue_number=int(r.get("gitea_issue_number", 0)),
linked_pr_numbers=pr_nums,
last_seen=r.get("last_seen"),
event_count=int(r.get("event_count", 1)),
created_at=r.get("created_at"),
updated_at=r.get("updated_at"),
)
)
bridge_active = any(p.bridge_enabled for p in providers)
return ObservabilitySnapshot(
schema_version=OBSERVABILITY_SCHEMA_VERSION,
providers=providers,
links=links,
total_links=len(links),
sentry_links_count=sentry_cnt,
glitchtip_links_count=glitchtip_cnt,
bridge_active=bridge_active,
)
def snapshot_to_dict(snapshot: ObservabilitySnapshot) -> dict[str, Any]:
return snapshot.to_dict()
-143
View File
@@ -1,143 +0,0 @@
"""HTML view renderer for the Sentry/GlitchTip observability console (#649, Phase 4).
Renders connection status widgets, error correlation links, and gated issue creation
affordances over the read-only observability snapshot.
"""
from __future__ import annotations
import html
from typing import Any
from webui.layout import render_page
from webui.observability_loader import ObservabilitySnapshot, snapshot_to_dict
def _badge(status: str) -> str:
st = (status or "").lower()
if st == "healthy":
return '<span class="badge badge-success">healthy</span>'
if st == "disabled":
return '<span class="badge badge-warning">disabled (dry-run)</span>'
if st in {"missing_token", "not_configured"}:
return f'<span class="badge badge-muted">{html.escape(st)}</span>'
return f'<span class="badge">{html.escape(st)}</span>'
def _provider_card(p: dict[str, Any]) -> str:
name = html.escape(str(p.get("provider", "provider")).upper())
base_url = html.escape(str(p.get("base_url", "")))
org = html.escape(str(p.get("org", "")))
proj = html.escape(str(p.get("project", "")))
status_badge = _badge(str(p.get("status", "")))
min_events = p.get("min_events_for_issue", 2)
lookback = html.escape(str(p.get("lookback", "24h")))
bridge_enabled = "yes" if p.get("bridge_enabled") else "no"
return f"""
<div class="card" style="margin-bottom: 1rem; padding: 1rem; border: 1px solid #ccc; border-radius: 6px;">
<div style="display: flex; justify-content: space-between; align-items: center;">
<h3 style="margin: 0;">{name} Connection</h3>
<div>{status_badge}</div>
</div>
<table style="width: 100%; margin-top: 0.5rem; border-collapse: collapse;">
<tr><td><strong>Base URL:</strong></td><td><code>{base_url}</code></td></tr>
<tr><td><strong>Scope:</strong></td><td><code>{org} / {proj}</code></td></tr>
<tr><td><strong>Bridge Enabled:</strong></td><td><code>{bridge_enabled}</code></td></tr>
<tr><td><strong>Min Events for Issue:</strong></td><td><code>{min_events}</code></td></tr>
<tr><td><strong>Lookback Window:</strong></td><td><code>{lookback}</code></td></tr>
</table>
</div>
"""
def render_observability_page(snapshot: ObservabilitySnapshot | dict[str, Any]) -> str:
"""Render the observability dashboard HTML page."""
data = snapshot.to_dict() if isinstance(snapshot, ObservabilitySnapshot) else dict(snapshot)
providers_raw = data.get("providers", [])
provider_cards = "".join(_provider_card(p) for p in providers_raw) if providers_raw else "<p>No providers configured.</p>"
links = data.get("links", [])
link_rows = []
for l in links:
prov = html.escape(str(l.get("provider", "")))
p_issue_id = html.escape(str(l.get("provider_issue_id", "")))
fingerprint = html.escape(str(l.get("fingerprint") or ""))
g_issue_num = int(l.get("gitea_issue_number", 0))
g_org = html.escape(str(l.get("gitea_org", "")))
g_repo = html.escape(str(l.get("gitea_repo", "")))
g_issue_link = f'<strong>#{g_issue_num}</strong> ({g_org}/{g_repo})'
event_cnt = int(l.get("event_count", 1))
last_seen = html.escape(str(l.get("last_seen") or ""))
short_id = html.escape(str(l.get("provider_short_id") or p_issue_id))
link_rows.append(f"""
<tr>
<td><code>{prov}</code></td>
<td><strong>{short_id}</strong><br><small style="color: #666;">id: {p_issue_id}</small></td>
<td><code>{fingerprint}</code></td>
<td>{g_issue_link}</td>
<td>{event_cnt}</td>
<td><small>{last_seen}</small></td>
</tr>
""")
table_body = "".join(link_rows) if link_rows else '<tr><td colspan="6" style="text-align: center; padding: 1.5rem; color: #666;">No correlated incident links stored. Bridge operates under dry-run default.</td></tr>'
metrics = data.get("metrics", {})
total_links = metrics.get("total_links", 0)
sentry_cnt = metrics.get("sentry_links_count", 0)
glitchtip_cnt = metrics.get("glitchtip_links_count", 0)
body_html = f"""
<h2>Observability & Incident Bridge (#649)</h2>
<p>Read-only console surface for Sentry/GlitchTip provider connections, error correlation,
and durable Gitea issue linkage.</p>
<div class="alert alert-info" style="background: #f0f4f8; padding: 1rem; border-left: 4px solid #0052cc; margin-bottom: 1.5rem;">
<strong>ADR Authority Model:</strong> Gitea records durable issue history. Control-plane DB coordinates incident links.
Sentry/GlitchTip observe errors. Raw monitoring incidents are <em>never</em> assignable control-plane work items.
Durable issue creation is gated and dry-runable via the <code>#612</code> bridge APIs.
</div>
<h3>Provider Connections</h3>
<div style="display: grid; grid-template-columns: repeat(auto-fit, minmax(300px, 1fr)); gap: 1rem; margin-bottom: 2rem;">
{provider_cards}
</div>
<div style="display: flex; justify-content: space-between; align-items: center; margin-bottom: 1rem;">
<h3 style="margin: 0;">Correlated Incidents ({total_links})</h3>
<div>
<span class="badge" style="margin-right: 0.5rem;">Sentry: {sentry_cnt}</span>
<span class="badge">GlitchTip: {glitchtip_cnt}</span>
</div>
</div>
<table class="table" style="width: 100%; border-collapse: collapse; border: 1px solid #ddd;">
<thead>
<tr style="background: #f9f9f9; text-align: left;">
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Provider</th>
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Incident ID</th>
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Fingerprint</th>
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Gitea Issue Link</th>
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Events</th>
<th style="padding: 0.5rem; border-bottom: 2px solid #ddd;">Last Seen</th>
</tr>
</thead>
<tbody>
{table_body}
</tbody>
</table>
<div style="margin-top: 2rem; padding: 1rem; background: #fafafa; border: 1px solid #eee; border-radius: 4px;">
<h4 style="margin-top: 0;">Reconcile & Link Controls (Gated)</h4>
<p style="margin-bottom: 0.5rem; color: #555;">
Create or reconcile durable Gitea issues from provider observations using the <code>#612</code> incident bridge:
</p>
<code>mcp call gitea_observability_reconcile_incident --provider sentry --apply false</code>
</div>
"""
return render_page(title="Observability", body_html=body_html)