Compare commits

...
Author SHA1 Message Date
sysadmin f0c6255d7d Merge remote-tracking branch 'prgs/master' into feat/issue-659-maintenance-drain-mode 2026-07-25 18:40:57 -04:00
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
jcwalker3 71031c812e Merge branch 'master' into feat/issue-648-notifications-console 2026-07-25 17:29:57 -05:00
jcwalker3 e43ddd3cbe docs(incident): retroactive audit for direct-to-master commit 2fa97c26 (#670) 2026-07-25 17:29:34 -05:00
sysadmin a64ba08e27 fix(webui): address #905 REQUEST_CHANGES on notifications classifier
B1: classify_attention_event uses structured flags/category only — never
substring-match human-authored title/summary for escalation.

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

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

Regression tests cover all three blockers.

Refs #648
2026-07-25 18:26:46 -04:00
sysadminandClaude Opus 4.8 77d808e7d4 merge(master): resolve PR #907 allocator drain vs side_effect_free
Keep #659 maintenance-drain assignment stop and #643 side_effect_free
apply guard; session registration remains gated by side_effect_free.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 18:12:03 -04:00
sysadminandClaude Opus 4.8 bb8c3a537b merge(master): resolve PR #905 conflicts with requests/linkage
Keep notifications (#648) routes and nav alongside master requests (#643)
and other base updates.

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

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 17:02:12 -04:00
sysadmin 4f06d30e07 feat(webui): implement notifications and human-attention routing (#648) 2026-07-25 16:41:16 -04:00
14 changed files with 2354 additions and 1 deletions
+41
View File
@@ -26,6 +26,7 @@ from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Mapping, Sequence
import maintenance_drain
from control_plane_db import (
ControlPlaneDB,
ControlPlaneError,
@@ -936,6 +937,46 @@ def allocate_next_work(
"allocation_mode": (allocation_mode or "").strip() or None,
}
# #659 AC2: while maintenance drain is active, no new work is assigned —
# for dry-run and apply alike, so a preview can never be read as evidence
# that work was assignable during the drain. Checked before session
# registration so a drained allocator leaves no new state behind.
try:
drain_record = db.read_maintenance_drain(remote=remote, org=org, repo=repo)
except Exception as exc: # noqa: BLE001 — unreadable drain state fails closed
return {
"success": False,
"outcome": OUTCOME_NO_SAFE,
"reasons": [
f"maintenance-drain state lookup failed: {exc} (fail closed, #659)"
],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
}
drain_decision = maintenance_drain.classify_assignment(drain_record)
if not drain_decision["assignment_allowed"]:
return {
"success": True,
"outcome": OUTCOME_WAIT,
"apply": apply,
"role": role_norm,
"allocation_mode": mode,
"remote": remote,
"org": org,
"repo": repo,
"selected": None,
"reasons": list(drain_decision["reasons"]),
"reason_code": drain_decision["reason_code"],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
"maintenance_drain": maintenance_drain.status_payload(
drain_record, remote=remote, org=org, repo=repo
),
}
# A side-effect-free run may never reserve: reserving is a write, and the
# flag is the caller's assertion that this call writes nothing (#643).
if side_effect_free and apply:
+194 -1
View File
@@ -31,8 +31,9 @@ from typing import Any, Iterator, Sequence
import dependency_graph
import gitea_audit
import maintenance_drain
SCHEMA_VERSION = 5
SCHEMA_VERSION = 6
# Assignable work kinds only — raw monitoring incidents are never work items.
WORK_KINDS = frozenset({"issue", "pr"})
@@ -239,6 +240,31 @@ CREATE INDEX IF NOT EXISTS idx_session_checkpoints_session
CREATE INDEX IF NOT EXISTS idx_session_checkpoints_work
ON session_checkpoints(remote, org, repo, work_kind, work_number);
-- Graceful maintenance-drain state (#659). One current row per repository
-- scope — drain is a *state*, not a history, so entering and exiting update
-- the same row and every transition is audited to ``events``. Creating the
-- table is the v5->v6 migration: additive, idempotent, and it never touches
-- prior tables. ``state`` is CHECK-constrained so an unknown value can never
-- be written and later read as "not draining".
CREATE TABLE IF NOT EXISTS maintenance_drain (
drain_id TEXT PRIMARY KEY,
remote TEXT NOT NULL,
org TEXT NOT NULL,
repo TEXT NOT NULL,
state TEXT NOT NULL DEFAULT 'inactive'
CHECK (state IN ('inactive', 'draining')),
reason TEXT NOT NULL DEFAULT '',
requested_by TEXT NOT NULL DEFAULT '',
requested_by_profile TEXT NOT NULL DEFAULT '',
session_id TEXT NOT NULL DEFAULT '',
entered_at TEXT NOT NULL DEFAULT '',
exited_at TEXT NOT NULL DEFAULT '',
drain_schema_version INTEGER NOT NULL DEFAULT 6,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
UNIQUE (remote, org, repo)
);
-- Model usage, token cost, latency, and performance events (#651)
CREATE TABLE IF NOT EXISTS usage_events (
usage_id INTEGER PRIMARY KEY AUTOINCREMENT,
@@ -3025,3 +3051,170 @@ class ControlPlaneDB:
"live_lease_id": None if live_lease_id is None else str(live_lease_id),
"reconcile_action": "reconcile_required" if stale else "safe_to_resume",
}
# ── Maintenance drain (#659) ─────────────────────────────────────────────
@staticmethod
def _maintenance_drain_row(row: sqlite3.Row | None) -> dict[str, Any] | None:
"""Convert a ``maintenance_drain`` row to a plain record."""
if row is None:
return None
return {key: row[key] for key in row.keys()}
def read_maintenance_drain(
self, *, remote: str, org: str, repo: str
) -> dict[str, Any] | None:
"""Return the current drain record for a scope, or None if never set.
None and a stored ``inactive`` row mean the same thing to callers —
``maintenance_drain.is_draining`` treats both as not draining — so the
read never has to invent a record to answer the gate.
"""
with self._tx(immediate=False) as conn:
row = conn.execute(
"""
SELECT * FROM maintenance_drain
WHERE remote = ? AND org = ? AND repo = ?
""",
(str(remote or ""), str(org or ""), str(repo or "")),
).fetchone()
return self._maintenance_drain_row(row)
def set_maintenance_drain(
self,
*,
remote: str,
org: str,
repo: str,
state: str,
reason: str = "",
requested_by: str = "",
requested_by_profile: str = "",
session_id: str = "",
) -> dict[str, Any]:
"""Enter or exit maintenance drain for one repository scope (AC1).
The state transition is audited to ``events`` — entering and exiting
are exactly the moments an operator has to be able to reconstruct
later. Re-entering an already-draining scope is idempotent: it refreshes
the reason/owner metadata, keeps the original ``entered_at``, and
records no duplicate transition event.
Capability authorization happens above this layer (the drain tasks
carry a non-``gitea.*`` permission in the task capability map); the DB
records who asked and why, and never grants the right itself.
"""
state_norm = maintenance_drain.normalize_state(state)
raw = {
"reason": str(reason or ""),
"requested_by": str(requested_by or ""),
"requested_by_profile": str(requested_by_profile or ""),
"session_id": str(session_id or ""),
}
clean = gitea_audit.redact(raw)
remote_s, org_s, repo_s = str(remote or ""), str(org or ""), str(repo or "")
now_s = _ts()
with self._tx() as conn:
existing = conn.execute(
"""
SELECT * FROM maintenance_drain
WHERE remote = ? AND org = ? AND repo = ?
""",
(remote_s, org_s, repo_s),
).fetchone()
prior_state = (
maintenance_drain.normalize_state(existing["state"])
if existing is not None
else maintenance_drain.STATE_INACTIVE
)
transitioned = prior_state != state_norm
prior_entered = (
str(existing["entered_at"] or "") if existing is not None else ""
)
prior_exited = (
str(existing["exited_at"] or "") if existing is not None else ""
)
if state_norm == maintenance_drain.STATE_DRAINING:
# A re-entry keeps the original entry time (the drain never
# stopped); a fresh entry stamps now and clears the old exit.
entered_at = prior_entered if (not transitioned and prior_entered) else now_s
exited_at = ""
else:
entered_at = prior_entered
exited_at = now_s if (transitioned or not prior_exited) else prior_exited
if existing is None:
drain_id = uuid.uuid4().hex
conn.execute(
"""
INSERT INTO maintenance_drain(
drain_id, remote, org, repo, state, reason,
requested_by, requested_by_profile, session_id,
entered_at, exited_at, drain_schema_version,
created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
drain_id, remote_s, org_s, repo_s, state_norm,
clean["reason"], clean["requested_by"],
clean["requested_by_profile"], clean["session_id"],
entered_at, exited_at,
maintenance_drain.DRAIN_SCHEMA_VERSION, now_s, now_s,
),
)
else:
drain_id = str(existing["drain_id"])
conn.execute(
"""
UPDATE maintenance_drain
SET state = ?, reason = ?, requested_by = ?,
requested_by_profile = ?, session_id = ?,
entered_at = ?, exited_at = ?,
drain_schema_version = ?, updated_at = ?
WHERE drain_id = ?
""",
(
state_norm, clean["reason"], clean["requested_by"],
clean["requested_by_profile"], clean["session_id"],
entered_at, exited_at,
maintenance_drain.DRAIN_SCHEMA_VERSION, now_s, drain_id,
),
)
if transitioned:
event_type = (
"maintenance_drain_enter"
if state_norm == maintenance_drain.STATE_DRAINING
else "maintenance_drain_exit"
)
conn.execute(
"""
INSERT INTO events(work_item_id, event_type, message, created_at)
VALUES (NULL, ?, ?, ?)
""",
(
event_type,
f"drain {drain_id} scope {remote_s}/{org_s}/{repo_s} "
f"{prior_state} -> {state_norm} by "
f"{clean['requested_by'] or '(unknown)'} "
f"({clean['requested_by_profile'] or 'no profile'}); "
f"reason: {clean['reason'] or '(none)'}",
now_s,
),
)
row = conn.execute(
"SELECT * FROM maintenance_drain WHERE drain_id = ?", (drain_id,)
).fetchone()
record = self._maintenance_drain_row(row) or {}
return {
"record": record,
"drain_id": drain_id,
"state": state_norm,
"prior_state": prior_state,
"transitioned": transitioned,
}
@@ -0,0 +1,83 @@
# Incident #670: bare direct-to-master commit `2fa97c26` (retroactive audit)
Status: verified; disposition recommendation: **accept as-is, no revert** (final
disposition owned by controller per issue #670).
## Summary
Commit `2fa97c26fbda555a1a83930ca5fdcea9d8e47b50`
(`fix(mcp): load dotenv relative to project root`) landed on `prgs/master`
as a single-parent commit with no PR wrapper and no review record, bypassing
the sanctioned issue → branch → PR → review → merge workflow. It was
discovered during the PR #654 post-merge audit. PR #654 itself merged
cleanly via the Gitea API and did **not** introduce this commit.
## Verification evidence (acceptance criteria 13)
- **AC1 — present on `prgs/master`: yes.**
`git merge-base --is-ancestor 2fa97c26fbda555a1a83930ca5fdcea9d8e47b50 prgs/master` → true.
- **AC2 — no PR or review record: confirmed.**
The commit is a single-parent, non-merge commit sitting directly on
first-parent master between the #629 merge (`5ab5fe85`) and the #654
merge (`ec903b0d`). A PR landing on master produces a merge commit (or a
PR-linked head); neither exists here. The controller audit at issue-create
time also found no PR wrapper and no review record for this SHA.
- **AC3 — changed files and diff summary: confirmed.**
`gitea_auth.py | 5 +++--` (+3/2). Single parent
`5ab5fe8583c07134d55dadf09381aecb67df246e`. The change moves
`PROJECT_ROOT` derivation above `load_dotenv()` and loads
`.env` relative to the project root instead of the process CWD.
## AC4 — why no immediate revert
- The dotenv fix is intentional and required for correct runtime behavior:
without it, `load_dotenv()` resolves `.env` against the process working
directory, which breaks MCP server launches whose CWD is not the project
root.
- The change is small (+3/2), self-contained in `gitea_auth.py`, and has
been running on master without incident since 2026-07-10.
- Reverting would re-introduce a real bug to remove a provenance defect —
the wrong trade. Provenance is repaired retroactively by this document,
issue #670, and the hardening landed under #671.
- If the controller later judges the change unsafe, a separate
revert/repair issue is the sanctioned path (issue #670, recommended
disposition option 4).
## AC5 — workflow-hardening linkage
Prevention already landed: **issue #671** (closed)
*“Block direct pushes to stable branches from MCP workflow sessions”*,
implemented by commit `5933d87647656643a67a50331c4c7b06ea751dad`
(`feat(guard): block direct stable-branch pushes from MCP workflow sessions`).
Shipped guardrails include:
- `gitea_record_stable_branch_push_attempt` — classifies proposed commands
for direct stable-branch push intent (`git push <remote> master`,
refspecs, `HEAD:master`, `--force`, dry-run intent, `:master` delete),
plus root/control-checkout local commits not carried by an issue branch,
and writes a durable `stable_branch_contamination` marker.
- `gitea_audit_stable_branch_contamination` — reconciler-only audit/clear
path; a contaminated worker session cannot self-clear.
- Review/merge/close/completion mutations fail closed while a
contamination marker is active.
## AC6 — PR #654 was not the source
- `2fa97c26` is the **first parent** of the #654 merge commit
`ec903b0d619e7a27d24aed272a890f4e5d381411`; it predates the #654 merge.
- First-parent history `5ab5fe8..ec903b0`:
`2fa97c2 fix(mcp): load dotenv relative to project root` followed by
`ec903b0 Merge pull request 'feat: lifecycle role/hazard labels ... (#603)' (#654)`.
- The #654 merger audit confirmed `ec903b0d` was a valid Gitea-API merge,
the `git push prgs master` attempt during that run was a no-op, and the
net change `2fa97c2..ec903b0` contained only the reviewed #603
lifecycle-label files.
- Conclusion: #654 merged reviewed content only; the unauthorized-path
defect is solely the earlier bare commit `2fa97c26`.
## Explicit non-actions (unchanged by this audit)
- No revert of `2fa97c26`.
- No force-push or history rewrite.
- No master mutation from the audit session.
+45
View File
@@ -0,0 +1,45 @@
# MCP maintenance-drain mode (#659)
Graceful **maintenance drain** stops new work assignment and defers non-allowlisted
mutations so sessions can finish critical handoffs and checkpoint before a
restart. It is **not** a restart authorization: the drain *proof* and apply gate
remain #661.
## State
Per repository scope (`remote`/`org`/`repo`) in the control-plane DB table
`maintenance_drain` (schema v6):
| State | Meaning |
|-------|---------|
| `inactive` | Normal operation (also: no row) |
| `draining` | Assignment stopped; non-allowlisted mutations deferred |
Enter/exit transitions are audited as `maintenance_drain_enter` /
`maintenance_drain_exit` events.
## Tools
| Tool | Permission | Effect |
|------|------------|--------|
| `gitea_maintenance_drain_status` | `gitea.read` | Observe drain (every session) |
| `gitea_enter_maintenance_drain` | `runtime.maintenance_drain` | Enter drain (capability-gated) |
| `gitea_exit_maintenance_drain` | `runtime.maintenance_drain` | Exit drain |
`runtime.maintenance_drain` is intentionally **not** a `gitea.*` op, so ordinary
author profiles cannot enter drain by accident.
## Enforcement
1. **Allocator** (`allocate_next_work`): while draining, returns `outcome=wait`
with `reason_code=maintenance_drain_assignment_stopped` for dry-run and apply.
2. **Mutation preflight** (`verify_preflight_purity`): non-allowlisted mutation
tasks raise `MaintenanceDrainError` with a typed next action.
3. **Allowlist** (safety only): heartbeats, lease release/abandon, session
checkpoints, enter/exit drain. Reads always work.
## Restart relationship
Drain mode prepares the blast radius. Restart apply still requires a clean
`DrainProof` (#661) or authorized break-glass. Status payloads never claim
restart permission.
+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.
+240
View File
@@ -1472,6 +1472,9 @@ def verify_preflight_purity(
# contaminated by manual MCP daemon process killing (reconciler-exempt).
_enforce_runtime_recovery_contamination_gate(task, remote)
# #659 AC3: defer non-allowlisted mutations while maintenance drain is active.
_enforce_maintenance_drain_gate(task, remote=remote, org=org, repo=repo)
ctx = _resolve_namespace_mutation_context(worktree_path)
workspace = ctx["workspace_path"]
canonical_root = ctx["canonical_repo_root"]
@@ -2004,6 +2007,55 @@ def _enforce_runtime_recovery_contamination_gate(
)
def _enforce_maintenance_drain_gate(
task: str | None,
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> None:
"""#659 AC3: defer non-allowlisted mutations while drain is active.
The single mutation chokepoint already used by every gated task, so drain
coverage cannot drift per-tool. Allowlisted safety operations (heartbeat,
release/abandon, checkpoint, drain exit) pass through so an in-flight
session can still finish and hand off; everything else is deferred with a
typed blocker. Unreadable drain state fails closed a drain that cannot be
read is not evidence that no drain is running.
"""
if _preflight_in_test_mode() and not os.environ.get(
"GITEA_TEST_FORCE_MAINTENANCE_DRAIN"
):
return
if maintenance_drain.is_allowlisted_task(task):
return
try:
_h, o, r = _resolve(remote, None, org, repo)
except Exception: # noqa: BLE001 — scope resolution is best-effort here
o, r = (org or ""), (repo or "")
db, errs = _control_plane_db_or_error()
if db is None:
raise RuntimeError(
"maintenance-drain state could not be read: "
f"{'; '.join(errs) or 'control-plane DB unavailable'} (fail closed, #659)"
)
try:
record = db.read_maintenance_drain(remote=remote or "", org=o, repo=r)
except Exception as exc: # noqa: BLE001
raise RuntimeError(
f"maintenance-drain state could not be read: {_redact(str(exc))} "
"(fail closed, #659)"
) from exc
decision = maintenance_drain.classify_mutation(task, record)
if not decision["allowed"]:
raise maintenance_drain.MaintenanceDrainError(
maintenance_drain.format_drain_block_error(decision),
decision=decision,
)
def _enforce_stable_branch_contamination_gate(
task: str | None,
remote: str | None = None,
@@ -2064,6 +2116,7 @@ import allocator_service # noqa: E402
import allocator_dependencies # noqa: E402
import dependency_graph # noqa: E402 # #784 durable dependency edges
import control_plane_db # noqa: E402
import maintenance_drain # noqa: E402 # #659 graceful maintenance-drain mode
import lease_lifecycle # noqa: E402
import lease_policy # noqa: E402
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
@@ -22559,6 +22612,193 @@ def gitea_workflow_dashboard(
return payload
@mcp.tool()
def gitea_maintenance_drain_status(
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Read-only: current maintenance-drain state for a repository scope (#659 AC4).
Every session must be able to observe drain so it can stop creating new work
and finish only allowlisted safety operations. Never mutates; never restarts.
"""
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"read_only": True,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
try:
_h, o, r = _resolve(remote, host, org, repo)
except ValueError as exc:
return {"success": False, "read_only": True, "reasons": [str(exc)]}
db, errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"read_only": True,
"reasons": errs or ["control-plane DB unavailable"],
"maintenance_drain": maintenance_drain.status_payload(
None, remote=remote, org=o, repo=r
),
}
try:
record = db.read_maintenance_drain(remote=remote, org=o, repo=r)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"read_only": True,
"reasons": [f"drain state unreadable: {_redact(str(exc))}"],
}
payload = maintenance_drain.status_payload(
record, remote=remote, org=o, repo=r
)
return {"success": True, "read_only": True, "maintenance_drain": payload}
@mcp.tool()
def gitea_enter_maintenance_drain(
reason: str = "",
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
session_id: str | None = None,
) -> dict:
"""Enter graceful maintenance-drain mode for a repository scope (#659 AC1).
Stops new assignment and defers non-allowlisted mutations until exit. Requires
``runtime.maintenance_drain`` (controller/lifecycle capability not granted
by ordinary Gitea author profiles). Audited in the control-plane event log.
"""
cap_block = _profile_operation_gate("runtime.maintenance_drain")
if cap_block:
return {
"success": False,
"performed": False,
"reasons": cap_block,
"permission_report": _permission_block_report(
"runtime.maintenance_drain"
),
}
try:
_h, o, r = _resolve(remote, host, org, repo)
except ValueError as exc:
return {"success": False, "performed": False, "reasons": [str(exc)]}
db, errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"performed": False,
"reasons": errs or ["control-plane DB unavailable"],
}
profile = get_profile() or {}
try:
result = db.set_maintenance_drain(
remote=remote,
org=o,
repo=r,
state=maintenance_drain.STATE_DRAINING,
reason=reason or "operator-entered maintenance drain",
requested_by=str(
(profile.get("identity") or {}).get("username")
or profile.get("expected_username")
or ""
),
requested_by_profile=str(profile.get("profile_name") or ""),
session_id=str(session_id or ""),
)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"performed": False,
"reasons": [f"enter drain failed: {_redact(str(exc))}"],
}
record = result.get("record") or {}
return {
"success": True,
"performed": True,
"transitioned": bool(result.get("transitioned")),
"state": result.get("state"),
"prior_state": result.get("prior_state"),
"drain_id": result.get("drain_id"),
"maintenance_drain": maintenance_drain.status_payload(
record, remote=remote, org=o, repo=r
),
}
@mcp.tool()
def gitea_exit_maintenance_drain(
reason: str = "",
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
session_id: str | None = None,
) -> dict:
"""Exit graceful maintenance-drain mode (#659 AC1). Restores assignment and mutations."""
cap_block = _profile_operation_gate("runtime.maintenance_drain")
if cap_block:
return {
"success": False,
"performed": False,
"reasons": cap_block,
"permission_report": _permission_block_report(
"runtime.maintenance_drain"
),
}
try:
_h, o, r = _resolve(remote, host, org, repo)
except ValueError as exc:
return {"success": False, "performed": False, "reasons": [str(exc)]}
db, errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"performed": False,
"reasons": errs or ["control-plane DB unavailable"],
}
profile = get_profile() or {}
try:
result = db.set_maintenance_drain(
remote=remote,
org=o,
repo=r,
state=maintenance_drain.STATE_INACTIVE,
reason=reason or "operator-exited maintenance drain",
requested_by=str(
(profile.get("identity") or {}).get("username")
or profile.get("expected_username")
or ""
),
requested_by_profile=str(profile.get("profile_name") or ""),
session_id=str(session_id or ""),
)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"performed": False,
"reasons": [f"exit drain failed: {_redact(str(exc))}"],
}
record = result.get("record") or {}
return {
"success": True,
"performed": True,
"transitioned": bool(result.get("transitioned")),
"state": result.get("state"),
"prior_state": result.get("prior_state"),
"drain_id": result.get("drain_id"),
"maintenance_drain": maintenance_drain.status_payload(
record, remote=remote, org=o, repo=r
),
}
@mcp.tool()
def gitea_request_mcp_restart(
remote: str = "dadeschools",
+281
View File
@@ -0,0 +1,281 @@
"""Graceful MCP maintenance-drain mode (#659).
Drain is the visible, capability-gated state that lets an operator stop new
work and quiesce mutations *before* a restart, instead of cutting sessions off
mid-mutation. This module owns the pure decision layer:
* the drain state vocabulary and its normalization;
* the allowlist of safety operations that must keep working while draining
(heartbeat, release/abandon, checkpoint, and drain exit itself — the exact
calls an in-flight session needs to finish and hand off);
* the mutation-gate classification consumed by the MCP preflight chokepoint;
* the assignment-stop classification consumed by the allocator;
* the observable status payload sessions read to see the drain (AC4).
Durable state lives in the control-plane DB (``maintenance_drain`` table);
enforcement lives at the existing chokepoints. Nothing here performs I/O, so
both callers can share one decision without importing each other.
Scope note: the machine-verifiable *drain proof* and the restart gate that
consumes it are #661's scope, not this module's. Drain here stops assignment
and mutation and makes the state observable; it never authorizes a restart.
"""
from __future__ import annotations
from typing import Any, Mapping
# ── State vocabulary ──────────────────────────────────────────────────────────
STATE_INACTIVE = "inactive"
STATE_DRAINING = "draining"
DRAIN_STATES = frozenset({STATE_INACTIVE, STATE_DRAINING})
# Typed blocker code surfaced to clients (never a bare string at call sites).
BLOCKER_DRAIN_ACTIVE = "maintenance_drain_active"
# Reason code for the allocator's assignment stop.
REASON_ASSIGNMENT_STOPPED = "maintenance_drain_assignment_stopped"
DRAIN_SCHEMA_VERSION = 6
class MaintenanceDrainError(RuntimeError):
"""Raised when a mutation is refused because drain is active (fail closed)."""
def __init__(self, message: str, *, decision: Mapping[str, Any] | None = None):
super().__init__(message)
self.decision = dict(decision or {})
self.reason_code = BLOCKER_DRAIN_ACTIVE
# ── Safety allowlist ──────────────────────────────────────────────────────────
# Mutations that stay permitted while draining. Every entry is a *quiesce*
# operation: it either proves an in-flight task is still alive, hands its claim
# back, records the durable state a restart needs, or ends the drain. Nothing
# that creates new work, new branches, new PRs, or new review/merge verdicts is
# on this list — that is the whole point of the drain.
ALLOWLISTED_DRAIN_TASKS: frozenset[str] = frozenset(
{
# Liveness of work already in flight.
"heartbeat_issue_lock",
"heartbeat_reviewer_pr_lease",
"post_heartbeat",
# Handing claims back so nothing is stranded across the restart.
"release_workflow_lease",
"release_reviewer_pr_lease",
"release_merger_pr_lease",
"abandon_workflow_lease",
# Durable recovery state (#660) must be writable *during* drain.
"write_session_checkpoint",
"checkpoint_session",
# The drain controls themselves — exit must never be self-blocked.
"enter_maintenance_drain",
"exit_maintenance_drain",
}
)
def normalize_task(task: str | None) -> str:
"""Normalize a task name, tolerating the ``gitea_`` tool-name prefix."""
name = str(task or "").strip()
if name.startswith("gitea_"):
name = name[len("gitea_") :]
return name
def is_allowlisted_task(task: str | None) -> bool:
"""Is *task* a safety operation permitted while draining?"""
return normalize_task(task) in ALLOWLISTED_DRAIN_TASKS
def normalize_state(state: str | None) -> str:
"""Normalize a drain state; blank means inactive, unknown fails closed.
Blank normalizes to ``inactive`` (no drain record = not draining), but an
unrecognized non-blank value raises: silently treating ``"drainig"`` as
inactive would disable the gate.
"""
value = str(state or "").strip().lower()
if not value:
return STATE_INACTIVE
if value not in DRAIN_STATES:
raise MaintenanceDrainError(
f"unknown maintenance-drain state {value!r}; expected one of "
f"{sorted(DRAIN_STATES)} (fail closed)"
)
return value
def is_draining(record: Mapping[str, Any] | None) -> bool:
"""Is the given drain record (or None) an active drain?"""
if not record:
return False
return normalize_state(record.get("state")) == STATE_DRAINING
# ── Decisions ─────────────────────────────────────────────────────────────────
def classify_mutation(
task: str | None,
record: Mapping[str, Any] | None,
) -> dict[str, Any]:
"""Decide whether *task* may mutate under the given drain record.
Returns a decision dict with ``allowed``/``deferred`` and, when refused, a
typed ``reason_code`` plus the one exact next action the caller may take.
Deferred (not failed): the operation is legal again after drain exits, so
the caller is told to wait rather than to retry a different way.
"""
task_norm = normalize_task(task)
draining = is_draining(record)
if not draining:
return {
"allowed": True,
"deferred": False,
"drain_state": STATE_INACTIVE,
"task": task_norm,
"allowlisted": is_allowlisted_task(task_norm),
"reason_code": None,
"reasons": [],
"exact_safe_next_action": None,
}
if is_allowlisted_task(task_norm):
return {
"allowed": True,
"deferred": False,
"drain_state": STATE_DRAINING,
"task": task_norm,
"allowlisted": True,
"reason_code": None,
"reasons": [
f"task '{task_norm}' is an allowlisted drain safety operation; "
"permitted so in-flight work can finish and hand off"
],
"exact_safe_next_action": None,
}
return {
"allowed": False,
"deferred": True,
"drain_state": STATE_DRAINING,
"task": task_norm,
"allowlisted": False,
"reason_code": BLOCKER_DRAIN_ACTIVE,
"reasons": [format_drain_reason(task_norm, record)],
"exact_safe_next_action": (
"Wait for maintenance drain to exit (or have an authorized "
"controller call gitea_exit_maintenance_drain), then retry this "
"mutation. Reads and gitea_maintenance_drain_status stay available."
),
}
def classify_assignment(record: Mapping[str, Any] | None) -> dict[str, Any]:
"""Decide whether the allocator may assign new work (AC2)."""
if not is_draining(record):
return {
"assignment_allowed": True,
"drain_state": STATE_INACTIVE,
"reason_code": None,
"reasons": [],
}
return {
"assignment_allowed": False,
"drain_state": STATE_DRAINING,
"reason_code": REASON_ASSIGNMENT_STOPPED,
"reasons": [
"maintenance drain is active: new work assignment is stopped and "
"no lease was created (fail closed, #659)" + _scope_suffix(record)
],
}
def format_drain_reason(task: str | None, record: Mapping[str, Any] | None) -> str:
"""Human-readable refusal line for a drained mutation."""
task_norm = normalize_task(task) or "(unnamed task)"
return (
f"maintenance drain is active: mutation '{task_norm}' is deferred; only "
"allowlisted drain safety operations "
f"({', '.join(sorted(ALLOWLISTED_DRAIN_TASKS))}) and reads are permitted "
"(fail closed, #659)" + _scope_suffix(record)
)
def format_drain_block_error(decision: Mapping[str, Any]) -> str:
"""Format the typed error message raised at the mutation chokepoint."""
reasons = list(decision.get("reasons") or [])
head = reasons[0] if reasons else "maintenance drain is active (fail closed)"
action = decision.get("exact_safe_next_action")
return f"{head}. Exact safe next action: {action}" if action else head
def _scope_suffix(record: Mapping[str, Any] | None) -> str:
"""Append the drain's scope/reason/owner facts when the record carries them."""
if not record:
return ""
bits: list[str] = []
scope = "/".join(
str(record.get(key) or "") for key in ("remote", "org", "repo")
).strip("/")
if scope:
bits.append(f"scope {scope}")
if record.get("reason"):
bits.append(f"reason: {record['reason']}")
if record.get("requested_by"):
bits.append(f"entered by {record['requested_by']}")
if record.get("entered_at"):
bits.append(f"at {record['entered_at']}")
return f" ({'; '.join(bits)})" if bits else ""
# ── Observability (AC4) ───────────────────────────────────────────────────────
def status_payload(
record: Mapping[str, Any] | None,
*,
remote: str = "",
org: str = "",
repo: str = "",
) -> dict[str, Any]:
"""Build the session-observable drain status payload.
Always answers, including when no drain record exists: an absent record is
a definitive "not draining", not an unknown.
"""
draining = is_draining(record)
rec: Mapping[str, Any] = record or {}
return {
"drain_state": STATE_DRAINING if draining else STATE_INACTIVE,
"draining": draining,
"remote": str(rec.get("remote") or "") or remote,
"org": str(rec.get("org") or "") or org,
"repo": str(rec.get("repo") or "") or repo,
"reason": str(rec.get("reason") or ""),
"requested_by": str(rec.get("requested_by") or ""),
"requested_by_profile": str(rec.get("requested_by_profile") or ""),
"session_id": str(rec.get("session_id") or ""),
"entered_at": str(rec.get("entered_at") or ""),
"exited_at": str(rec.get("exited_at") or ""),
"assignment_stopped": draining,
"mutations_deferred": draining,
"allowlisted_tasks": sorted(ALLOWLISTED_DRAIN_TASKS),
"reads_permitted": True,
"record_present": bool(record),
"schema_version": DRAIN_SCHEMA_VERSION,
"drain_proof_scope": (
"drain proof and the restart gate that consumes it are #661 scope; "
"this status never authorizes a restart"
),
"safe_next_action": (
"Wait for drain to exit before retrying deferred mutations; "
"allowlisted safety operations and reads remain available."
if draining
else "None; maintenance drain is not active."
),
}
+18
View File
@@ -414,6 +414,24 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"role": "controller",
},
# #659 maintenance drain. Same reasoning as the lifecycle controls above:
# entering/exiting drain quiesces a whole namespace, so it carries a
# non-``gitea.*`` permission that no configured Gitea profile satisfies by
# accident (AC1 — capability-gated and audited). Reading drain state is
# ordinary read authority: every session must be able to see the drain (AC4).
"enter_maintenance_drain": {
"permission": "runtime.maintenance_drain",
"role": "controller",
},
"exit_maintenance_drain": {
"permission": "runtime.maintenance_drain",
"role": "controller",
},
"maintenance_drain_status": {
"permission": "gitea.read",
"role": "author",
},
# #601 first-class lease lifecycle — inspect/list need read; mutations gate on
# ownership in the control-plane DB (not a separate Gitea write permission).
"list_workflow_leases": {
+237
View File
@@ -0,0 +1,237 @@
"""Tests for graceful MCP maintenance-drain mode (#659).
Acceptance coverage:
1. Enter/exit is durable and audited (DB substrate).
2. New work assignment stops during drain (allocator WAIT).
3. Mutations deferred except allowlisted safety ops.
4. Sessions can observe drain state.
5. Fail-closed on unreadable drain state.
"""
from __future__ import annotations
import sys
import tempfile
import unittest
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
import maintenance_drain
from control_plane_db import ControlPlaneDB
from allocator_service import WorkCandidate, allocate_next_work, OUTCOME_WAIT
class TestDrainDecisions(unittest.TestCase):
def test_inactive_allows_mutations_and_assignment(self):
decision = maintenance_drain.classify_mutation("create_pr", None)
self.assertTrue(decision["allowed"])
self.assertFalse(decision["deferred"])
assign = maintenance_drain.classify_assignment(None)
self.assertTrue(assign["assignment_allowed"])
def test_draining_defers_non_allowlisted_mutation(self):
record = {"state": "draining", "remote": "prgs", "org": "o", "repo": "r"}
decision = maintenance_drain.classify_mutation("create_pr", record)
self.assertFalse(decision["allowed"])
self.assertTrue(decision["deferred"])
self.assertEqual(decision["reason_code"], maintenance_drain.BLOCKER_DRAIN_ACTIVE)
self.assertIn("create_pr", decision["reasons"][0])
def test_allowlisted_safety_ops_pass_during_drain(self):
record = {"state": "draining"}
for task in (
"heartbeat_issue_lock",
"gitea_release_reviewer_pr_lease",
"write_session_checkpoint",
"exit_maintenance_drain",
):
with self.subTest(task=task):
decision = maintenance_drain.classify_mutation(task, record)
self.assertTrue(decision["allowed"], decision)
def test_assignment_stopped_during_drain(self):
record = {"state": "draining", "reason": "upgrade"}
decision = maintenance_drain.classify_assignment(record)
self.assertFalse(decision["assignment_allowed"])
self.assertEqual(
decision["reason_code"], maintenance_drain.REASON_ASSIGNMENT_STOPPED
)
def test_unknown_state_fails_closed(self):
with self.assertRaises(maintenance_drain.MaintenanceDrainError):
maintenance_drain.normalize_state("drainig")
def test_status_payload_always_answers(self):
inactive = maintenance_drain.status_payload(None, remote="prgs", org="o", repo="r")
self.assertFalse(inactive["draining"])
self.assertTrue(inactive["reads_permitted"])
active = maintenance_drain.status_payload(
{"state": "draining", "reason": "reboot", "requested_by": "ops"},
remote="prgs",
org="o",
repo="r",
)
self.assertTrue(active["draining"])
self.assertTrue(active["assignment_stopped"])
self.assertTrue(active["mutations_deferred"])
self.assertIn("heartbeat_issue_lock", active["allowlisted_tasks"])
class TestDrainDB(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.db = ControlPlaneDB(db_path=str(Path(self._tmp.name) / "cp.sqlite3"))
def test_enter_exit_idempotent_and_audited(self):
first = self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="draining",
reason="planned restart",
requested_by="sysadmin",
requested_by_profile="prgs-controller",
session_id="s1",
)
self.assertTrue(first["transitioned"])
self.assertEqual(first["state"], "draining")
self.assertTrue(maintenance_drain.is_draining(first["record"]))
again = self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="draining",
reason="still draining",
requested_by="sysadmin",
requested_by_profile="prgs-controller",
session_id="s1",
)
self.assertFalse(again["transitioned"])
self.assertEqual(again["record"]["entered_at"], first["record"]["entered_at"])
exited = self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="inactive",
reason="done",
requested_by="sysadmin",
requested_by_profile="prgs-controller",
session_id="s1",
)
self.assertTrue(exited["transitioned"])
self.assertFalse(maintenance_drain.is_draining(exited["record"]))
self.assertTrue(exited["record"]["exited_at"])
# Events recorded for transitions only (enter + exit).
with self.db._tx(immediate=False) as conn:
rows = conn.execute(
"SELECT event_type FROM events WHERE event_type LIKE 'maintenance_drain_%' "
"ORDER BY event_id"
).fetchall()
types = [r[0] for r in rows]
self.assertEqual(types, ["maintenance_drain_enter", "maintenance_drain_exit"])
def test_read_missing_is_none_not_error(self):
self.assertIsNone(
self.db.read_maintenance_drain(remote="prgs", org="o", repo="r")
)
class TestAllocatorStopsDuringDrain(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.db = ControlPlaneDB(db_path=str(Path(self._tmp.name) / "cp.sqlite3"))
def test_allocate_returns_wait_while_draining(self):
self.db.set_maintenance_drain(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
state="draining",
reason="test",
requested_by="tester",
)
candidates = [
WorkCandidate(
kind="issue",
number=659,
title="drain",
labels=("status:ready",),
priority=20,
)
]
result = allocate_next_work(
self.db,
role="author",
session_id="test-session",
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
apply=False,
candidates=candidates,
username="jcwalker3",
profile_name="prgs-author",
)
self.assertEqual(result["outcome"], OUTCOME_WAIT)
self.assertIsNone(result.get("selected"))
self.assertEqual(
result.get("reason_code"),
maintenance_drain.REASON_ASSIGNMENT_STOPPED,
)
self.assertTrue(result["maintenance_drain"]["draining"])
def test_allocate_works_when_inactive(self):
candidates = [
WorkCandidate(
kind="issue",
number=659,
title="drain",
labels=("status:ready",),
priority=20,
)
]
result = allocate_next_work(
self.db,
role="author",
session_id="test-session-2",
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
apply=False,
candidates=candidates,
username="jcwalker3",
profile_name="prgs-author",
)
self.assertNotEqual(
result.get("reason_code"),
maintenance_drain.REASON_ASSIGNMENT_STOPPED,
)
class TestCapabilityMap(unittest.TestCase):
def test_drain_tasks_mapped(self):
import task_capability_map as tcm
self.assertEqual(
tcm.required_permission("enter_maintenance_drain"),
"runtime.maintenance_drain",
)
self.assertEqual(
tcm.required_permission("exit_maintenance_drain"),
"runtime.maintenance_drain",
)
self.assertEqual(
tcm.required_permission("maintenance_drain_status"),
"gitea.read",
)
if __name__ == "__main__":
unittest.main()
+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)
+24
View File
@@ -80,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
@@ -889,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.
@@ -1020,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"]),
+1
View File
@@ -46,6 +46,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
NavItem("/queue", "Queue"),
NavItem("/leases", "Leases"),
NavItem("/actions", "Actions"),
NavItem("/notifications", "Notifications"),
NavItem("/requests", "Requests"),
)),
NavGroup("Runtime/Sessions", (
+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()