Compare commits

..
Author SHA1 Message Date
sysadminandClaude Opus 5 97bc190fc2 docs(remote-mcp): inventory stdio- and localhost-coupled assumptions (#930)
Add docs/remote-mcp/coupling-inventory.md, the blocking first child of epic
#929. It enumerates every place gitea_mcp_server.py and its supporting modules
depend on being a local, client-spawned, stdio-attached process on the
operator's machine.

62 entries across the seven required categories: transport bind, launch
provenance, role binding, credentials, runtime freshness, local filesystem,
and durable state. Each entry carries a file and line anchor resolving at
7bf4f12584, states what the code assumes today
and what it would observe on a remote host, is classified as portable as
written / needs a seam / needs a replacement / cannot be remote, and is
assigned to exactly one epic child. Every child from #931 through #939 is
named by at least one entry. Summary tables count entries per category, per
classification, per category-by-classification, and per child.

Documentation only. No server behavior changes.

Closes #930

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01HvJz7bUz5CkZgUxq8twHMz
2026-07-26 02:09:09 -04:00
sysadmin 7bf4f12584 Merge pull request 'fix(guardrail): detect and reject manually launched duplicate MCP role servers (#686)' (#925) from fix/issue-686-detect-reject-manual-mcp into master 2026-07-25 18:32:23 -05:00
sysadmin dc0bff764d test(runtime): patch _trusted_session_repository in test_activate_profile_succeeds_when_enabled 2026-07-25 19:27:12 -04:00
sysadmin dfe8d7c28d fix(mcp): detect and reject manually launched duplicate MCP role servers (Closes #686) 2026-07-25 19:25:23 -04:00
sysadmin 2b4e43042a Merge pull request 'feat(tests): add concurrent-session MCP restart safety tests (Closes #666)' (#910) from feat/issue-666-concurrent-mcp-restart-tests into master 2026-07-25 17:44:09 -05:00
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
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
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
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
23 changed files with 3907 additions and 15 deletions
+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.
+12
View File
@@ -153,7 +153,19 @@ not a tool argument: a session must never be able to authorize itself.
## Related
- #630 — manual daemon killing as contaminated recovery (this contrast, enforced).
- #657 — restart-path inventory and daemon classification.
- #686 — manual server launch detection & fail-closed provenance gate.
- #531 / #544 — stale-runtime detection (`ps`-based); sibling failure mode.
- #558 / `docs/mcp-daemon-import-guard.md` — why shell imports are not a repair.
- `docs/mcp-client-registration.md` — per-server registration contract.
- `docs/mcp-namespace-health.md` — probe sources and mutation enforcement.
## Sanctioned reconnect vs forbidden manual launch (#686)
In addition to manual process killing (#630), manually launching a duplicate role server from an ad hoc shell (`python3 mcp_server.py`) is forbidden and fail-closed:
- **Why manual launches are unsupported:** A terminal-launched `mcp_server.py` holds its own stdio transport; it can never bind to the IDE client's stdio pipes. It cannot restore a dropped IDE namespace, and a manual duplicate process masks stale client-managed runtimes for that profile, defeating stale-runtime gates.
- **Sanctioned path:** Supported recovery is IDE/client-managed reconnect only (`/mcp reconnect`, IDE restart, or sanctioned reconnect exposure).
- **Fail-closed enforcement (#686):** Mutating tools on a server lacking client-managed launch provenance (`GITEA_CLIENT_MANAGED=1`) refuse execution fail-closed with typed blocker `unsupported_manual_launch` and an exact next action. Unsupported `GITEA_*` env overrides (e.g. `GITEA_DUMMY`) are surfaced in diagnostics rather than silently ignored.
- **Inventory & staleness:** Staleness diagnostics ignore non-client-managed duplicates when evaluating runtime freshness and inventory duplicate processes per profile (#657, #686).
+230
View File
@@ -0,0 +1,230 @@
# Remote-MCP coupling inventory
Every place the Gitea MCP server depends on being a local, client-spawned, stdio-attached
process on the operator's machine.
- **Issue:** #930 (Remote-MCP 01), child 1 of epic #929.
- **Generated against commit:** `7bf4f1258451823a55b36d2157e74f8457165088` (`master`).
- **Anchors:** every `file:line` below resolves at the commit above and at the commit that
adds this document. This change adds one new file and edits no existing file, so no
existing line number shifts between the two.
- **Scope:** documentation only. No server behavior changes in this child.
## How to read an entry
| Field | Meaning |
| ----- | ------- |
| **Anchor** | `file:line` at the commit under review. |
| **Assumes today** | What the code takes for granted while running as a local stdio process. |
| **Observes remotely** | What the same code would actually see on a shared remote host. |
| **Class** | One of: *portable as written*, *needs a seam*, *needs a replacement*, *cannot be remote*. |
| **Owner** | Exactly one epic child (#931#939) responsible for the fix. |
Classification meanings:
- **portable as written** — the code is already transport-, host-, and principal-neutral; it
moves unchanged once its inputs are supplied by a remote-aware caller.
- **needs a seam** — the logic is correct but is wired to a hard-coded local source. It needs
an injection point, not new semantics.
- **needs a replacement** — the semantics themselves are local-only. A remote deployment
needs a differently-defined mechanism, not the same mechanism relocated.
- **cannot be remote** — the operation is inherently about the operator's own machine
(its process table, its keychain, its checkout). It must either stay local behind an
explicit boundary or be deleted from the remote surface.
---
## 1. Transport bind
The transport is bound literally, once, at process start, and the bound value is the root of
the mutation-authorization chain.
| ID | Anchor | Assumes today | Observes remotely | Class | Owner |
| -- | ------ | ------------- | ----------------- | ----- | ----- |
| T1 | `gitea_mcp_server.py:23750` | The single production bind call passes the literal `transport="stdio"` immediately before the server loop. | The literal is wrong for any non-stdio deployment; there is no parameter to change it. | needs a seam | #931 |
| T2 | `mcp_daemon_guard.py:45` | `_PRODUCTION_TRANSPORTS = frozenset({"stdio"})` is the closed allowlist of production transports. | A remote transport name is rejected by the allowlist before any other check runs. | needs a seam | #931 |
| T3 | `mcp_daemon_guard.py:174` | `bind_native_mcp_transport` raises `UnsanctionedRuntimeError` for any transport outside `_PRODUCTION_TRANSPORTS` (raise at `mcp_daemon_guard.py:187`). | The remote server fails to start rather than degrading; the failure is correct, but the allowlist is the only thing that must change. | needs a seam | #931 |
| T4 | `mcp_daemon_guard.py:328` | `is_native_mcp_transport()` asserts a process-local runtime record whose `pid` matches `os.getpid()` and whose phase is `transport_bound`. The predicate itself names no transport. | Unchanged semantics: one server process that bound one transport. It stays true on a remote host. | portable as written | #931 |
| T5 | `mcp_daemon_guard.py:349` | `is_production_native_mcp_transport()` adds only a `mode == production` check on top of T4. | Unchanged. | portable as written | #931 |
| T6 | `irrecoverable_provenance.py:497` | `assess_transport_for_auth_mint()` requires production native transport before minting non-forgeable recovery authorization (#709 F1). | The gate is transport-agnostic in form, but its guarantee — "an ordinary Python process cannot reach this" — is currently underwritten by the stdio bind. Under a remote transport the guarantee must be re-derived from the authenticated session, not from the bind. | needs a seam | #931 |
| T7 | `gitea_mcp_server.py:8375` | Consumer: refuses to proceed unless `assess_transport_for_auth_mint()` allows. | Unchanged given a corrected T6. | portable as written | #931 |
| T8 | `gitea_mcp_server.py:8624` | Second consumer of the same gate on the confirmation path. | Unchanged given a corrected T6. | portable as written | #931 |
| T9 | `mcp_server.py:4` | Module docstring asserts "Runs over stdio." as a property of the server. | The stated contract becomes false on the remote deployment and is load-bearing documentation for operators. | needs a replacement | #931 |
## 2. Launch provenance
Mutations fail closed unless the process can prove a client launched it with real stdio pipes
and `GITEA_CLIENT_MANAGED` provenance. Every proof in this section is a statement about the
local operating system.
| ID | Anchor | Assumes today | Observes remotely | Class | Owner |
| -- | ------ | ------------- | ----------------- | ----- | ----- |
| P1 | `gitea_mcp_server.py:14588` | `_is_client_managed_process()` derives provenance from `GITEA_CLIENT_MANAGED` / `GITEA_MCP_CLIENT_MANAGED` / `GITEA_SERVER_PROVENANCE` / `GITEA_FORCE_CLIENT_MANAGED` on this process's own environment. | A long-lived remote process has one environment for all callers, so a per-process env var can no longer say anything about the caller that issued a request. | needs a replacement | #934 |
| P2 | `gitea_mcp_server.py:14606` | Falls back to `sys.stdin.isatty()`: an active TTY on stdin means a human launched it from a terminal, so refuse. | A remote server has no meaningful stdin. The signal is absent, not merely different. | cannot be remote | #934 |
| P3 | `gitea_mcp_server.py:14618` | `_provenance_mutation_block()` emits `blocker_kind: "unsupported_manual_launch"` and a "reconnect the IDE/client-managed MCP namespace" remediation. | The block shape is reusable; its predicate and its remediation text are both stdio-specific. | needs a seam | #934 |
| P4 | `gitea_mcp_server.py:20599` | `_check_mcp_runtimes_diagnostics()` shells `ps -o pid,lstart,command -ax` and greps for `mcp_server.py` to find peer role servers. | On a shared host the process table lists unrelated tenants' processes, or none at all under a container. Peer discovery by `ps` has no remote meaning. | cannot be remote | #934 |
| P5 | `gitea_mcp_server.py:20702` | More than one process per `GITEA_MCP_PROFILE` in the local process table is reported as a duplicate-launch fault. | A remote endpoint is expected to serve many concurrent sessions per role. "Two processes for one role" becomes the normal case, so the check inverts from a safety net into a false wall. | cannot be remote | #934 |
| P6 | `gitea_mcp_server.py:20715` | Processes lacking client-managed provenance are ignored for runtime freshness and reported as manual launches. | Same defect as P5: correctness depends on enumerating local peers. | cannot be remote | #934 |
| P7 | `gitea_config.py:1172` | `RECOGNIZED_GITEA_ENV_KEYS` is the allowlist of `GITEA_*` env vars a legitimately launched server may carry; anything else is contamination. | Configuration on a remote host arrives from deployment tooling, not from a client-authored env block. The allowlist keeps working mechanically but stops proving anything about provenance. | needs a replacement | #934 |
| P8 | `gitea_mcp_server.py:20683` | The unsupported-env scan applies `RECOGNIZED_GITEA_ENV_KEYS` to *other* processes' environments harvested via `ps eww <pid>`. | Reading another process's environment is unavailable or prohibited across tenants, and is not exposed in this form outside macOS/BSD `ps`. | cannot be remote | #934 |
| P9 | `mcp_daemon_guard.py:126` | `mark_sanctioned_daemon()` requires the claiming stack frame's resolved absolute path to be the canonical `mcp_server.py` / `gitea_mcp_server.py` next to the guard module; basename spoofing is rejected. | Entrypoint-path identity still exists on a remote host, but it authenticates the *deployment*, not the *caller*. It must be kept and demoted from "authorizes mutations" to "authorizes the process". | needs a seam | #934 |
| P10 | `gitea_config.py:1233` | The client-config generator emits `"GITEA_CLIENT_MANAGED": "1"` into each generated MCP client entry, alongside `GITEA_MCP_CONFIG` / `GITEA_MCP_PROFILE`. | A remote endpoint is addressed by URL and credential, not by a spawn command with an env block. This generator produces the wrong artifact entirely. | needs a replacement | #938 |
| P11 | `mcp_namespace_health.py:232` | Namespace health classifies a namespace as `client_managed` or `manual_launch` from the reported env summary. | During dual-run, local and remote namespaces coexist and must both be classifiable; a two-valued local/manual axis cannot express "remote endpoint, authenticated session". | needs a replacement | #939 |
| P12 | `gitea_mcp_server.py:18161` | The diagnostics payload reports `server_provenance` as exactly `"client_managed"` or `"manual_launch"`. | This is the field a cutover operator reads to confirm which deployment served a call. It must gain a remote value before dual-run parity can be validated. | needs a replacement | #939 |
## 3. Role binding
Role separation is currently enforced by *which process a call reaches*. The process is pinned
to one role for its lifetime by an environment variable.
| ID | Anchor | Assumes today | Observes remotely | Class | Owner |
| -- | ------ | ------------- | ----------------- | ----- | ----- |
| R1 | `gitea_config.py:54` | `ENV_PROFILE = "GITEA_MCP_PROFILE"` is the single source of the active profile, read from the process environment. | One shared process serves several principals; a process-wide profile cannot answer "who is calling now". This is the root of the coupling. | needs a replacement | #932 |
| R2 | `review_workflow_load.py:95` | Reads `GITEA_MCP_PROFILE` directly to decide the reviewer workflow binding. | Reads the deployment's profile, not the caller's, silently granting or denying the wrong role. | needs a replacement | #932 |
| R3 | `mcp_discoverability.py:152` | Reads `GITEA_MCP_PROFILE` to describe the namespace to the client. | Correct logic, wrong input source; it needs the request principal injected. | needs a seam | #932 |
| R4 | `webui/deployment_boundary.py:115` | Reads `GITEA_MCP_PROFILE` to classify the deployment boundary for the console. | Same as R3. | needs a seam | #932 |
| R5 | `gitea_mcp_server.py:21106` | Remediation text instructs the operator to "Relaunch the server with `GITEA_MCP_PROFILE` set to a profile that has the required permission". | Relaunching a shared remote endpoint to change one caller's role is not a valid instruction; it would re-role every other session. | needs a replacement | #932 |
| R6 | `native_mcp_preference.py:93` | Detects shell commands that override `GITEA_MCP_PROFILE` away from the session (`native_mcp_preference.py:223`) and flags them as CLI auth divergence. | The divergence check is genuinely useful and survives, but its notion of "the session's profile" must come from the request principal. | needs a seam | #932 |
| R7 | `gitea_mcp_server.py:20671` | Recovers a peer server's role by regexing `GITEA_MCP_PROFILE=` out of that process's environment. | Depends on P4/P8 process-table access; role discovery by peer-env scraping has no remote analogue. | cannot be remote | #932 |
## 4. Credentials
Every token resolves, directly or indirectly, from one human's macOS keychain.
| ID | Anchor | Assumes today | Observes remotely | Class | Owner |
| -- | ------ | ------------- | ----------------- | ----- | ----- |
| C1 | `gitea_config.py:956` | `_keychain_token()` shells `security find-generic-password -s <item> -w`. | `security(1)` is a macOS binary reading the calling user's login keychain. It does not exist on a Linux host and would be the wrong identity even on a shared Mac. | cannot be remote | #933 |
| C2 | `gitea_config.py:974` | `resolve_token(profile, keychain_lookup=_keychain_token)` dispatches on `auth.type` of `env` or `keychain`, defaulting the lookup to C1. | The injectable `keychain_lookup` parameter is the existing seam; a remote credential provider plugs in here without changing the dispatch. | needs a seam | #933 |
| C3 | `gitea_config.py:1015` | `keychain_auth(item_id)` constructs the `{"type": "keychain", "id": ...}` reference stored in profiles. | The reference type itself encodes "macOS keychain" into persisted config. A remote provider needs a new auth reference type, not a new value of this one. | needs a replacement | #933 |
| C4 | `mcp_daemon_guard.py:440` | `assert_keychain_access_allowed()` fails closed for git-credential keychain fill outside a sanctioned daemon, with an operator opt-out env var. | The gate protects a mechanism that will not exist remotely. Its replacement must gate the *credential provider* call, not the keychain call, or the protection silently lapses. | needs a replacement | #933 |
| C5 | `sentry_incident_bridge.py:190` | `resolve_token(env)` resolves the Sentry token from an injected env mapping with no keychain path. | Already host-neutral; it is the shape the Gitea credential path should converge on. | portable as written | #933 |
| C6 | `gitea_mcp_server.py:18469` | The profile-audit tool calls `gitea_config.resolve_token(p)` for every configured profile to report "credentials present" without networking. | On a remote host this would materialize every principal's credential inside one process — an audit surface that becomes a credential-aggregation risk. | needs a seam | #933 |
## 5. Runtime freshness
The mutation gate is defined as "the commit this process started at matches the checkout on
this disk, and both match live master". Two of those three terms are local-disk facts.
| ID | Anchor | Assumes today | Observes remotely | Class | Owner |
| -- | ------ | ------------- | ----------------- | ----- | ----- |
| F1 | `master_parity_gate.py:168` | `capture_startup_parity(root)` reads git `HEAD` from the server's own root once at startup and returns it as the baseline. | A remote host carries a deployed artifact, not the operator's checkout. Its `HEAD` says nothing about the operator's working tree, which is the thing the gate exists to protect. | cannot be remote | #935 |
| F2 | `master_parity_gate.py:255` | `mutation_safe = determinable and in_parity and live_known and not live_stale` — a conjunction of two local-HEAD comparisons and one live-remote comparison. | Two of the three conjuncts lose meaning, so the whole verdict does. A remote deployment needs a redefined, testable freshness predicate rather than this one relocated. | needs a replacement | #935 |
| F3 | `master_parity_gate.py:164` | The live-remote head is probed and cached per `(root, remote, branch)`, keyed on the local root. | The live-remote probe is the one conjunct that survives; it needs a key that is not the operator's filesystem path. | needs a seam | #935 |
| F4 | `gitea_mcp_server.py:18262` | `gitea_assess_master_parity` publishes `startup_head` / `local_head` / `live_remote_head` / `mutation_safe` as the authoritative mutation-safety verdict. | The tool's contract is consumed by every mutation caller and by the operator; it must keep its shape while its semantics are redefined, or every consumer breaks at once. | needs a replacement | #935 |
| F5 | `gitea_mcp_server.py:23054` | Falls back to `_process_boot_head_sha` — the commit this process booted at — when the parity payload has no `startup_head`. | Same defect as F1, in a fallback path that is easy to miss when F1 is fixed. | needs a seam | #935 |
| F6 | `gitea_mcp_server.py:20615` | Staleness is also inferred from `os.path.getmtime()` of `gitea_mcp_server.py` under `PROJECT_ROOT` (`gitea_mcp_server.py:20611`), compared against peer process start times. | File mtime on a deployed artifact tracks the deploy, not the operator's edits, and the peer start times it is compared against come from the unavailable process table (P4). | cannot be remote | #935 |
## 6. Local filesystem
Author and reviewer tools act directly on the operator's checkout.
| ID | Anchor | Assumes today | Observes remotely | Class | Owner |
| -- | ------ | ------------- | ----------------- | ----- | ----- |
| L1 | `gitea_mcp_server.py:10122` | `gitea_bootstrap_author_issue_worktree` creates and binds a git worktree on the server's own disk. | The remote host has no operator checkout to add a worktree to. Executing this remotely would act on the wrong disk while reporting success. | cannot be remote | #936 |
| L2 | `gitea_mcp_server.py:190` | `ACTIVE_WORKTREE_ENV = "GITEA_ACTIVE_WORKTREE"` and `AUTHOR_WORKTREE_ENV` (`gitea_mcp_server.py:191`) carry the active workspace as process-wide environment. | Process-wide workspace state cannot represent per-session workspaces on a shared endpoint. | needs a replacement | #936 |
| L3 | `gitea_mcp_server.py:9801` | Binding a worktree writes `os.environ["GITEA_AUTHOR_WORKTREE"]` and `os.environ["GITEA_ACTIVE_WORKTREE"]` (`gitea_mcp_server.py:9802`), mutating global process state. | One session's bind would silently retarget every other concurrent session in the same process. This is a correctness bug the moment concurrency is real. | needs a replacement | #936 |
| L4 | `reviewer_inventory_worktree.py:48` | `_BRANCHES_WORKTREE_RE = re.compile(r"\bbranches/", re.I)` requires review worktree paths to sit under `branches/`. | A path convention on the operator's machine, asserted as a validation rule. It needs to become a property of a declared workspace, not a substring test. | needs a seam | #936 |
| L5 | `stable_control_runtime.py:54` | `DEV_WORKTREE_SEGMENT = "branches"` classifies a process root as a development worktree by path segment. | Same class of assumption as L4, on the runtime-classification side. | needs a seam | #936 |
| L6 | `mcp_server.py:42` | `check_conflict_markers()` runs at import and `os.walk`s the install directory for unresolved conflict markers, `sys.exit(1)` on a hit. | On a remote host it scans a deployed artifact, which by construction never has conflict markers — so the guard passes trivially and stops protecting the thing it was written to protect. | needs a replacement | #936 |
| L7 | `role_session_router.py:487` | `check_mid_merge()` reports infra-stop from `.git/MERGE_HEAD`, `rebase-merge`, `rebase-apply` and a source conflict scan under the server's project root. | Same inversion as L6: it would report the deployment's git state, not the operator's. | needs a replacement | #936 |
| L8 | `author_issue_bootstrap.py:996` | Enumerates worktrees with `git -C <root> worktree list --porcelain`. | Requires a real local clone with real worktrees; there is nothing equivalent to enumerate remotely. | cannot be remote | #936 |
| L9 | `mcp_server.py:10` | Redirects `sys.stderr` to the fixed path `/tmp/mcp_server_stderr.log` outside pytest. | A single fixed `/tmp` path is shared by every concurrent server on a host and is not a deployment's logging surface. | needs a replacement | #938 |
| L10 | `gitea_mcp_server.py:2314` | `ISSUE_LOCK_FILE = "/tmp/gitea_issue_lock.json"` — the legacy single global lock slot. | One global `/tmp` slot per host cannot represent concurrent remote sessions and is world-visible on a shared machine. | needs a replacement | #937 |
| L11 | `issue_lock_provenance.py:14` | `ISSUE_LOCK_FILE = os.environ.get("GITEA_ISSUE_LOCK_FILE", "/tmp/gitea_issue_lock.json")` keeps the same `/tmp` default in the provenance path. | Same as L10; the env override is a local escape hatch, not a remote design. | needs a replacement | #937 |
## 7. Durable state
Locks, leases, session state, and the control-plane database live in the operator's home
directory and are keyed on local PIDs.
| ID | Anchor | Assumes today | Observes remotely | Class | Owner |
| -- | ------ | ------------- | ----------------- | ----- | ----- |
| S1 | `issue_lock_store.py:26` | `DEFAULT_LOCK_DIR = ~/.cache/gitea-tools/issue-locks` — per-issue lock files under one user's home. | A shared endpoint has no single operator home; per-user paths make locks invisible across sessions and hosts. | needs a replacement | #937 |
| S2 | `issue_lock_store.py:83` | `session_pointer_path()` names the session pointer file `session-<os.getpid()>.json`. | Many sessions share one PID on a remote server, so the pointer collapses to a single slot and sessions overwrite each other. | cannot be remote | #937 |
| S3 | `issue_lock_store.py:98` | `is_process_alive(pid)` decides lock liveness by probing the local process table. | A PID recorded by one host is meaningless on another, and may coincidentally match a live unrelated process. | cannot be remote | #937 |
| S4 | `issue_lock_store.py:213` | Lock records stamp `session_pid` and `pid` from `os.getpid()`. | The recorded identity no longer distinguishes sessions; ownership checks silently pass for the wrong caller. | needs a replacement | #937 |
| S5 | `mcp_session_state.py:27` | `DEFAULT_STATE_DIR = ~/.cache/gitea-tools/session-state`, mode `0o700`. | Same home-directory coupling as S1, for review decision locks and workflow proofs. | needs a replacement | #937 |
| S6 | `mcp_session_state.py:559` | Session bodies stamp `session_pid` and `writer_pid` from `os.getpid()` (`mcp_session_state.py:560`). | Writer attribution collapses across concurrent sessions in one process. | needs a replacement | #937 |
| S7 | `control_plane_db.py:47` | `DEFAULT_DB_PATH = ~/.cache/gitea-tools/control-plane/control_plane.sqlite3`. | A per-user SQLite file is not reachable by, or safe for, multiple remote sessions or multiple hosts. | needs a replacement | #937 |
| S8 | `control_plane_db.py:386` | `sqlite3.connect(self.db_path, timeout=30)` — single-writer file locking tuned for one local process. | SQLite's write lock does not extend across hosts and degrades sharply under real concurrency; the store needs a concurrency-safe backend. | needs a replacement | #937 |
| S9 | `control_plane_db.py:1145` | Lease rows record `owner_pid` defaulting to `os.getpid()` (also `control_plane_db.py:2039`). | PID-keyed lease ownership is unusable across hosts and ambiguous within one shared process. | cannot be remote | #937 |
| S10 | `mcp_daemon_guard.py:53` | `_DEFAULT_SESSION_STATE_DIR` is pinned once at transport bind so a later `GITEA_MCP_SESSION_STATE_DIR` change cannot manufacture a second authority domain (#695 AC2). | The single-authority-domain invariant is exactly right and must be preserved; only its backing location needs to move. | needs a seam | #937 |
| S11 | `gitea_mcp_server.py:11875` | Reviewer-lease reclaim reads `owner_pid_alive` from the lease freshness record to decide whether an owner is dead. | Consumes S3/S9; a false "owner alive" or "owner dead" here reclaims or refuses a live lease. This is the highest-consequence consumer of PID liveness. | cannot be remote | #937 |
---
## Summary
### Entries per category
| Category | Entries |
| -------- | ------: |
| 1. Transport bind | 9 |
| 2. Launch provenance | 12 |
| 3. Role binding | 7 |
| 4. Credentials | 6 |
| 5. Runtime freshness | 6 |
| 6. Local filesystem | 11 |
| 7. Durable state | 11 |
| **Total** | **62** |
No category is empty, so no "this category has no coupling" justification is required.
### Entries per classification
| Classification | Entries |
| -------------- | ------: |
| portable as written | 5 |
| needs a seam | 16 |
| needs a replacement | 26 |
| cannot be remote | 15 |
| **Total** | **62** |
### Category × classification
| Category | portable | seam | replacement | cannot | Total |
| -------- | -------: | ---: | ----------: | -----: | ----: |
| 1. Transport bind | 4 | 4 | 1 | 0 | 9 |
| 2. Launch provenance | 0 | 2 | 5 | 5 | 12 |
| 3. Role binding | 0 | 3 | 3 | 1 | 7 |
| 4. Credentials | 1 | 2 | 2 | 1 | 6 |
| 5. Runtime freshness | 0 | 2 | 2 | 2 | 6 |
| 6. Local filesystem | 0 | 2 | 7 | 2 | 11 |
| 7. Durable state | 0 | 1 | 6 | 4 | 11 |
| **Total** | **5** | **16** | **26** | **15** | **62** |
### Entries per epic child
Every child from 2 through 10 is named by at least one entry, and every entry names exactly
one child.
| Child | Issue | Title | Entries | IDs |
| ----: | ----- | ----- | ------: | --- |
| 2 | #931 | Transport-neutral bind seam | 9 | T1T9 |
| 3 | #932 | Per-request principal resolution | 7 | R1R7 |
| 4 | #933 | Server-side credential provider | 6 | C1C6 |
| 5 | #934 | Remote-session provenance | 9 | P1P9 |
| 6 | #935 | Redefined master-parity gate | 6 | F1F6 |
| 7 | #936 | Local-filesystem vs remotable tool split | 8 | L1L8 |
| 8 | #937 | Concurrency-safe session, lock, and lease state | 13 | L10, L11, S1S11 |
| 9 | #938 | Authenticated remote MCP endpoint | 2 | P10, L9 |
| 10 | #939 | Dual-run cutover and rollback | 2 | P11, P12 |
| | | **Total** | **62** | |
## Notes for downstream children
- **The three highest-risk entries are P5, F2, and S11.** Each is a guard that does not
merely stop working remotely — it inverts. P5 turns concurrency into a reported fault,
F2 returns a verdict computed from terms that no longer mean anything, and S11 reclaims
or refuses leases on a PID-liveness answer that is wrong rather than unknown. A gate that
fails open while still reporting green is worse than one that fails to start.
- **T4, T5, T7, T8, and C5 are the portable core.** They show the target shape: predicates
over injected inputs, with no reference to the host, the process table, or the operator's
disk.
- **The keychain seam already exists** at C2 (`resolve_token`'s injectable `keychain_lookup`).
#933 should widen that seam rather than introduce a parallel path, and must remember C4 —
the guard protecting the old mechanism has to be re-pointed, or the protection lapses
silently when the mechanism is replaced.
- **`branches/` appears as a validation rule in at least two independent places** (L4, L5).
Path-substring conventions tend to have more copies than expected; #936 should re-grep
rather than trust this list to be exhaustive for that specific pattern.
+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.
+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.
+50 -1
View File
@@ -1169,10 +1169,57 @@ def server_command():
return python, [os.path.join(root, "mcp_server.py")]
RECOGNIZED_GITEA_ENV_KEYS = frozenset({
"GITEA_MCP_CONFIG",
"GITEA_MCP_PROFILE",
"GITEA_PROFILE_NAME",
"GITEA_SERVICE",
"GITEA_EXECUTION_ROLE",
"GITEA_CLIENT_MANAGED",
"GITEA_MCP_CLIENT_MANAGED",
"GITEA_SERVER_PROVENANCE",
"GITEA_AUTHOR_WORKTREE",
"GITEA_ACTIVE_WORKTREE",
"GITEA_DISABLE_KEYCHAIN",
"GITEA_CONTROL_PLANE_DB",
"GITEA_DB_PATH",
"GITEA_LOG_LEVEL",
"GITEA_DEBUG",
"GITEA_HMAC_SECRET",
"GITEA_IRRECOVERABLE_HMAC_SECRET",
"GITEA_FORCE_MCP_RUNTIME_CHECK",
"GITEA_FORCE_CLIENT_MANAGED",
})
RECOGNIZED_GITEA_ENV_PREFIXES = (
"GITEA_TOKEN_",
"GITEA_PASS_",
"GITEA_USER_",
"GITEA_URL_",
"GITEA_HOST_",
"GITEA_REMOTE_",
"GITEA_HTTP_HEADER_",
)
def get_unconsumed_gitea_env_overrides(env=None) -> dict[str, str]:
"""Find unsupported GITEA_* env vars present in *env* (defaults to os.environ)."""
target = os.environ if env is None else env
unconsumed = {}
for key, value in target.items():
if key.startswith("GITEA_"):
if key in RECOGNIZED_GITEA_ENV_KEYS:
continue
if any(key.startswith(p) for p in RECOGNIZED_GITEA_ENV_PREFIXES):
continue
unconsumed[key] = str(value)
return unconsumed
def launcher_entry(profile_name, config_path=None):
"""Return a thin MCP launcher entry for *profile_name*.
Contains only command/args and the two GITEA_MCP_* env vars — never a token
Contains command/args and the GITEA_MCP_* / GITEA_CLIENT_MANAGED env vars — never a token
or password. Suitable for Claude / Gemini / Codex ``mcpServers`` blocks.
"""
command, args = server_command()
@@ -1183,11 +1230,13 @@ def launcher_entry(profile_name, config_path=None):
"env": {
"GITEA_MCP_CONFIG": config_path or DEFAULT_CONFIG_PATH,
"GITEA_MCP_PROFILE": profile_name,
"GITEA_CLIENT_MANAGED": "1",
},
}
}
def keychain_set(item_id, token, account=None, runner=subprocess.run):
"""Store *token* in the macOS keychain under service *item_id*.
+120 -7
View File
@@ -14585,6 +14585,56 @@ def _session_context_mutation_block(
return blocked
def _is_client_managed_process() -> bool:
"""Check whether the current MCP server process has client-managed launch provenance (#686)."""
val = (
os.environ.get("GITEA_CLIENT_MANAGED")
or os.environ.get("GITEA_MCP_CLIENT_MANAGED")
or os.environ.get("GITEA_SERVER_PROVENANCE")
or os.environ.get("GITEA_FORCE_CLIENT_MANAGED")
or ""
).strip().lower()
if val in ("0", "false", "no", "manual", "manual_launch"):
return False
if val in ("1", "true", "yes", "client_managed"):
return True
# A terminal launch has an active TTY on stdin
try:
if sys.stdin and sys.stdin.isatty():
return False
except Exception:
pass
# Standard client launch or test runner with stdio pipe and profile env
if "GITEA_MCP_CONFIG" in os.environ or "GITEA_MCP_PROFILE" in os.environ or "GITEA_PROFILE_NAME" in os.environ:
return True
return False
def _provenance_mutation_block(**extra_fields) -> dict | None:
"""Refuse mutating tool calls on processes lacking client-managed launch provenance (#686)."""
if _is_client_managed_process():
return None
unconsumed = gitea_config.get_unconsumed_gitea_env_overrides()
blocked = {
"success": False,
"performed": False,
"blocker_kind": "unsupported_manual_launch",
"reasons": [
"mutation denied: server process was launched manually from a terminal without client-managed provenance (fail closed). Manually launched mcp_server.py processes cannot receive IDE stdio or serve workflow mutations."
],
"exact_next_action": "BLOCKED + RECONNECT: Reconnect the IDE/client-managed MCP server namespace instead of an ad hoc terminal launch. Hand-launched processes and mcp_config.json hand-edits are classified as workflow contamination.",
"provenance": "manual_launch",
"unconsumed_gitea_env": unconsumed,
}
blocked.update(extra_fields)
return blocked
def _profile_permission_block(required_operation: str, **extra_fields) -> dict | None:
"""Structured operation-gate denial for gated tools (#69, #142, #897).
@@ -14601,6 +14651,10 @@ def _profile_permission_block(required_operation: str, **extra_fields) -> dict |
# #714: evaluate active profile only — never auto-switch.
_ensure_matching_profile(required_operation, req_role, extra_fields.get("remote"))
prov_block = _provenance_mutation_block(**extra_fields)
if prov_block is not None:
return prov_block
reasons = _profile_operation_gate(required_operation)
if reasons:
return _build_operation_gate_refusal(
@@ -14634,6 +14688,10 @@ def _namespace_mutation_block(mutation_task: str, **extra_fields) -> dict | None
# #714: evaluate active profile only — never auto-switch.
_ensure_matching_profile(required_permission, required_role, extra_fields.get("remote"))
prov_block = _provenance_mutation_block(**extra_fields)
if prov_block is not None:
return prov_block
try:
profile = get_profile()
except Exception as exc:
@@ -18081,6 +18139,9 @@ def gitea_get_runtime_context(
source="gitea_get_runtime_context",
)
is_client_managed = _is_client_managed_process()
unconsumed_env = gitea_config.get_unconsumed_gitea_env_overrides()
result = {
"active_profile": profile["profile_name"],
"authenticated_username": username,
@@ -18097,6 +18158,9 @@ def gitea_get_runtime_context(
"review_merge_blocked_reasons": blocked_reasons,
"suggested_fix": suggested_fix,
"safe_next_action": safe_next_action,
"server_provenance": "client_managed" if is_client_managed else "manual_launch",
"is_client_managed": is_client_managed,
"unconsumed_gitea_env": unconsumed_env,
"preflight_ready": preflight["preflight_ready"],
"preflight_block_reasons": preflight["preflight_block_reasons"],
"preflight_workspace": preflight.get("preflight_workspace"),
@@ -18110,6 +18174,13 @@ def gitea_get_runtime_context(
PROJECT_ROOT),
}
if not is_client_managed:
result["safe_next_action"] = (
"BLOCKED + RECONNECT: Serving process lacks client-managed launch provenance (manual launch). "
"Reconnect the IDE/client-managed MCP server namespace instead of an ad hoc terminal launch."
)
# #702: read-only visibility into the inherited GITEA_ACTIVE_WORKTREE
# binding; recovery itself runs during capability resolution.
try:
@@ -20567,7 +20638,9 @@ def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) ->
self_pid = os.getpid()
self_stale = False
running_profiles = {}
all_profile_procs: dict[str, list[dict]] = {}
unsupported_env_found = set()
for line in proc.stdout.splitlines()[1:]:
line = line.strip()
if not line or "mcp_server.py" not in line:
@@ -20599,16 +20672,55 @@ def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) ->
if match:
profile = match.group(1)
is_client_managed = bool(
re.search(r'\bGITEA_CLIENT_MANAGED=(1|true|yes|client_managed)\b', env_out, re.IGNORECASE)
or re.search(r'\bGITEA_MCP_CLIENT_MANAGED=(1|true|yes|client_managed)\b', env_out, re.IGNORECASE)
or re.search(r'\bGITEA_SERVER_PROVENANCE=client_managed\b', env_out, re.IGNORECASE)
)
for env_match in re.finditer(r'\b(GITEA_[A-Z0-9_]+)=([^\s]+)', env_out):
k, v = env_match.group(1), env_match.group(2)
if k not in gitea_config.RECOGNIZED_GITEA_ENV_KEYS and not any(k.startswith(p) for p in gitea_config.RECOGNIZED_GITEA_ENV_PREFIXES):
unsupported_env_found.add(f"{k}={v}")
is_stale = (start_time < code_mtime) or git_stale
if pid == self_pid and is_stale:
self_stale = True
if profile not in running_profiles or start_time > running_profiles[profile]["start_time"]:
running_profiles[profile] = {
"pid": pid,
"start_time": start_time,
"is_stale": is_stale
}
proc_info = {
"pid": pid,
"start_time": start_time,
"is_stale": is_stale,
"is_client_managed": is_client_managed,
}
if profile not in all_profile_procs:
all_profile_procs[profile] = []
all_profile_procs[profile].append(proc_info)
running_profiles = {}
for profile, procs in all_profile_procs.items():
if len(procs) > 1:
pids_str = ", ".join(str(p["pid"]) for p in procs)
reasons.append(
f"stale-runtime: Duplicate MCP server process(es) detected for profile '{profile}' (PIDs: {pids_str}). "
"Manual or duplicate launches defeat staleness detection and cannot receive client stdio."
)
client_procs = [p for p in procs if p["is_client_managed"]]
if client_procs:
client_procs.sort(key=lambda p: p["start_time"], reverse=True)
running_profiles[profile] = client_procs[0]
else:
pids_str = ", ".join(str(p["pid"]) for p in procs)
reasons.append(
f"stale-runtime: Manually launched MCP process(es) detected without client-managed provenance for profile '{profile}' (PIDs: {pids_str}). "
"Manual launches cannot serve client stdio and are ignored for runtime freshness."
)
if unsupported_env_found:
reasons.append(
f"unsupported-env: Unsupported GITEA_* environment variable override(s) detected: {', '.join(sorted(unsupported_env_found))}. "
"Unknown env overrides are unsupported."
)
if self_stale:
# #685: report-only — no config utime, no thread, no os._exit.
@@ -20643,6 +20755,7 @@ def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) ->
return reasons
@mcp.tool()
def gitea_resolve_task_capability(
task: str,
+16
View File
@@ -225,6 +225,16 @@ def classify_namespace_probe(
# on bad data without treating success as IDE proof).
blocks = namespace_health_blocks_task("merge_pr", healthy)
import gitea_config
raw_env = process.get("env") if isinstance(process, dict) else None
unconsumed_env = gitea_config.get_unconsumed_gitea_env_overrides(raw_env)
is_client_managed = bool(
env_summary.get("GITEA_CLIENT_MANAGED") in ("1", "true", "yes", "client_managed")
or env_summary.get("GITEA_MCP_CLIENT_MANAGED") in ("1", "true", "yes", "client_managed")
or env_summary.get("GITEA_SERVER_PROVENANCE") == "client_managed"
)
provenance = "client_managed" if is_client_managed else "manual_launch"
return {
"success": healthy,
"healthy": healthy,
@@ -240,6 +250,9 @@ def classify_namespace_probe(
"error_message": error_message or None,
"reasons": reasons,
"remediation": remediation,
"provenance": provenance,
"is_client_managed": is_client_managed,
"unconsumed_gitea_env": unconsumed_env,
"diagnostics": {
"namespace": ns,
"required_tool": tool,
@@ -248,6 +261,9 @@ def classify_namespace_probe(
"env": env_summary,
"config_path": config_path,
"probe_source": source,
"provenance": provenance,
"is_client_managed": is_client_managed,
"unconsumed_gitea_env": unconsumed_env,
},
"blocks_merge_workflow": blocks,
}
+2
View File
@@ -44,6 +44,8 @@ def _reset_mutation_authority(monkeypatch):
]:
monkeypatch.delenv(env_key, raising=False)
monkeypatch.setenv("GITEA_CLIENT_MANAGED", "1")
# Isolate durable session-state files so tests never share host cache (#559).
import tempfile
+2 -2
View File
@@ -35,7 +35,7 @@ CONFIG = {
],
"forbidden_operations": [],
"execution_profile": "full-author",
"allowed_repositories": ["Example-Org/Example-Repo"],
"allowed_repositories": ["Scaled-Tech-Consulting/Gitea-Tools", "Example-Org/Example-Repo"],
},
"reviewer-no-commit": {
"enabled": True,
@@ -50,7 +50,7 @@ CONFIG = {
"gitea.repo.commit", "gitea.pr.create", "gitea.branch.push"
],
"execution_profile": "reviewer-no-commit",
"allowed_repositories": ["Example-Org/Example-Repo"],
"allowed_repositories": ["Scaled-Tech-Consulting/Gitea-Tools", "Example-Org/Example-Repo"],
},
},
"rules": {"allow_runtime_switching": False},
+1 -1
View File
@@ -175,7 +175,7 @@ class TestLauncherSnippets(unittest.TestCase):
def test_only_safe_keys_no_secrets(self):
entry = gitea_config.launcher_entry("prgs", "/cfg/profiles.json")["gitea-tools"]
self.assertEqual(set(entry), {"command", "args", "env"})
self.assertEqual(set(entry["env"]), {"GITEA_MCP_CONFIG", "GITEA_MCP_PROFILE"})
self.assertEqual(set(entry["env"]), {"GITEA_MCP_CONFIG", "GITEA_MCP_PROFILE", "GITEA_CLIENT_MANAGED"})
self.assertEqual(entry["env"]["GITEA_MCP_PROFILE"], "prgs")
blob = json.dumps(entry).lower()
for word in ("token", "password", "secret"):
@@ -0,0 +1,139 @@
"""Tests for Issue #686: Detect and reject manually launched duplicate MCP role servers."""
import os
import unittest
from unittest.mock import patch, MagicMock
from datetime import datetime
import gitea_config
import gitea_mcp_server
import mcp_namespace_health
class TestIssue686ManualMcpProvenance(unittest.TestCase):
def test_client_managed_process_detection(self):
"""Test _is_client_managed_process correctly detects provenance markers."""
with patch.dict(os.environ, {"GITEA_CLIENT_MANAGED": "1"}, clear=True):
self.assertTrue(gitea_mcp_server._is_client_managed_process())
with patch.dict(os.environ, {"GITEA_MCP_CLIENT_MANAGED": "true"}, clear=True):
self.assertTrue(gitea_mcp_server._is_client_managed_process())
with patch.dict(os.environ, {"GITEA_SERVER_PROVENANCE": "client_managed"}, clear=True):
self.assertTrue(gitea_mcp_server._is_client_managed_process())
with patch.dict(os.environ, {"GITEA_CLIENT_MANAGED": "0"}, clear=True):
self.assertFalse(gitea_mcp_server._is_client_managed_process())
def test_unconsumed_gitea_env_overrides(self):
"""Test surfacing of unsupported GITEA_* env overrides (e.g. GITEA_DUMMY)."""
env = {
"GITEA_MCP_PROFILE": "prgs-author",
"GITEA_CLIENT_MANAGED": "1",
"GITEA_DUMMY": "2",
"GITEA_UNKNOWN_FLAG": "abc",
}
unconsumed = gitea_config.get_unconsumed_gitea_env_overrides(env)
self.assertIn("GITEA_DUMMY", unconsumed)
self.assertEqual(unconsumed["GITEA_DUMMY"], "2")
self.assertIn("GITEA_UNKNOWN_FLAG", unconsumed)
self.assertNotIn("GITEA_MCP_PROFILE", unconsumed)
self.assertNotIn("GITEA_CLIENT_MANAGED", unconsumed)
def test_manual_server_mutation_fail_closed(self):
"""AC 2: Mutating tools on a server without client-managed provenance fail closed with a typed blocker."""
with patch.dict(os.environ, {"GITEA_CLIENT_MANAGED": "0"}, clear=True):
block = gitea_mcp_server._provenance_mutation_block(task="create_issue")
self.assertIsNotNone(block)
self.assertFalse(block["success"])
self.assertFalse(block["performed"])
self.assertEqual(block["blocker_kind"], "unsupported_manual_launch")
self.assertEqual(block["provenance"], "manual_launch")
self.assertTrue(any("mutation denied: server process was launched manually" in r for r in block["reasons"]))
self.assertIn("BLOCKED + RECONNECT", block["exact_next_action"])
def test_client_managed_server_mutation_passes_provenance_gate(self):
"""AC 3: Clean client-managed baseline passes the provenance gate."""
with patch.dict(os.environ, {"GITEA_CLIENT_MANAGED": "1"}, clear=True):
block = gitea_mcp_server._provenance_mutation_block(task="create_issue")
self.assertIsNone(block)
@patch("subprocess.run")
@patch("os.path.getmtime")
@patch("os.path.exists")
@patch("os.getpid")
def test_manual_duplicate_does_not_mask_stale_runtime(
self, mock_getpid, mock_exists, mock_getmtime, mock_run
):
"""AC 1 & AC 3: Staleness detection ignores manual duplicates and reports stale supported runtimes."""
mock_getpid.return_value = 12345
mock_exists.return_value = True
code_time = datetime(2026, 7, 8, 14, 0, 0)
mock_getmtime.return_value = code_time.timestamp()
# PID 12345: stale client-managed process (started at 13:00)
# PID 99999: fresh manual duplicate process (started at 15:00, no GITEA_CLIENT_MANAGED)
ps_output = (
" PID LSTART COMMAND\n"
"12345 Wed Jul 8 13:00:00 2026 /path/to/python mcp_server.py\n"
"99999 Wed Jul 8 15:00:00 2026 /path/to/python mcp_server.py\n"
)
mock_run_ps = MagicMock()
mock_run_ps.stdout = ps_output
mock_env_12345 = MagicMock()
mock_env_12345.stdout = "GITEA_MCP_PROFILE=prgs-author GITEA_CLIENT_MANAGED=1"
mock_env_99999 = MagicMock()
mock_env_99999.stdout = "GITEA_MCP_PROFILE=prgs-author GITEA_DUMMY=2"
def side_effect(args, **kwargs):
if args[0] == "ps" and "eww" in args:
pid = args[2]
if pid == "12345":
return mock_env_12345
elif pid == "99999":
return mock_env_99999
elif args[0] == "ps":
return mock_run_ps
raise ValueError(f"Unexpected args: {args}")
mock_run.side_effect = side_effect
reasons = gitea_mcp_server._check_mcp_runtimes_diagnostics("create_issue", ["prgs-author"])
# Manual duplicate process must be flagged
self.assertTrue(any("Duplicate MCP server process(es) detected" in r for r in reasons))
# Unsupported env override (GITEA_DUMMY=2) must be flagged
self.assertTrue(any("unsupported-env: Unsupported GITEA_* environment variable override(s) detected: GITEA_DUMMY=2" in r for r in reasons))
# Stale runtime must NOT be masked by fresh manual process 99999!
self.assertTrue(any("All matching profiles for task 'create_issue' (['prgs-author']) are running but stale" in r for r in reasons))
def test_namespace_health_classification_includes_provenance(self):
"""AC 1 & 4: mcp_namespace_health diagnostics include provenance and unconsumed_gitea_env."""
process = {
"pid": 5555,
"profile": "prgs-author",
"env": {
"GITEA_MCP_PROFILE": "prgs-author",
"GITEA_DUMMY": "99",
},
}
res = mcp_namespace_health.classify_namespace_probe(
"gitea-author",
configured=True,
registered_tools=["gitea_whoami"],
probe_result={"success": True},
process=process,
probe_source="client_namespace",
)
self.assertEqual(res["provenance"], "manual_launch")
self.assertFalse(res["is_client_managed"])
self.assertEqual(res["unconsumed_gitea_env"], {"GITEA_DUMMY": "99"})
self.assertEqual(res["diagnostics"]["provenance"], "manual_launch")
if __name__ == "__main__":
unittest.main()
+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()
+3 -3
View File
@@ -35,10 +35,10 @@ class TestMcpStaleRuntime(unittest.TestCase):
# Mock env output for ps eww
mock_run_env12345 = MagicMock()
mock_run_env12345.stdout = "GITEA_MCP_PROFILE=prgs-reconciler"
mock_run_env12345.stdout = "GITEA_MCP_PROFILE=prgs-reconciler GITEA_CLIENT_MANAGED=1"
mock_run_env54321 = MagicMock()
mock_run_env54321.stdout = "GITEA_MCP_PROFILE=prgs-author"
mock_run_env54321.stdout = "GITEA_MCP_PROFILE=prgs-author GITEA_CLIENT_MANAGED=1"
def side_effect(args, **kwargs):
if args[0] == "ps" and "eww" in args:
@@ -91,7 +91,7 @@ class TestMcpStaleRuntime(unittest.TestCase):
mock_run_ps.stdout = ps_output
mock_run_env = MagicMock()
mock_run_env.stdout = "GITEA_MCP_PROFILE=prgs-author"
mock_run_env.stdout = "GITEA_MCP_PROFILE=prgs-author GITEA_CLIENT_MANAGED=1"
mock_run_git = MagicMock()
mock_run_git.stdout = "FAKE2" # different SHA
+2 -1
View File
@@ -243,9 +243,10 @@ class TestRuntimeClarity(unittest.TestCase):
self.assertIn("switching is disabled", res["message"].lower())
self.assertIsNone(gitea_config._active_profile_override)
@patch("mcp_server._trusted_session_repository", return_value={"repository": "Example-Org/Example-Repo", "org": "Example-Org", "repo": "Example-Repo", "reasons": []})
@patch("mcp_server.api_request")
@patch("mcp_server.get_auth_header")
def test_activate_profile_succeeds_when_enabled(self, mock_auth, mock_api):
def test_activate_profile_succeeds_when_enabled(self, mock_auth, mock_api, mock_trusted):
self._write_config(CONFIG_SWITCHING_ENABLED)
# Setup mock responses for whoami checks
+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)
+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()
+61
View File
@@ -47,6 +47,9 @@ 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,
@@ -77,6 +80,11 @@ 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
@@ -416,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()
@@ -859,6 +894,22 @@ 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.
@@ -990,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"]),
@@ -1004,6 +1058,13 @@ 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"]),
+2
View File
@@ -46,10 +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", (
+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()
+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)
)