Compare commits

...
Author SHA1 Message Date
sysadmin 2b3f5baaeb fix(mcp): invalidate review session state on cross-profile activation (Closes #690)
A mid-run profile switch (reviewer -> author -> reviewer) left workflow-load
proof, reviewer lease binding, the review decision lock, live namespace
health, and preflight identity/capability stamps intact, so a formal verdict
could be recorded under contaminated session state.

- gitea_activate_profile now invalidates all review-critical session state
  on a cross-profile switch, in memory and in durable state keyed by either
  profile identity, and reports the invalidation + re-preflight requirement.
- Full reviewer preflight (whoami, load_review_workflow,
  resolve_task_capability(review_pr), head re-pin, lease re-acquire) is
  required before any formal verdict after a switch; switching back cannot
  resurrect the stale run.
- Namespace provenance: optional launcher-declared GITEA_MCP_NAMESPACE is
  reported by whoami/runtime context/capability resolution, and a declared
  namespace that disagrees with a task's required namespace fails closed.
- Docs: supported pattern is separate session/namespace per role, not
  in-process profile hopping mid-review.
2026-07-25 19:17:19 -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
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 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 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
14 changed files with 2282 additions and 1 deletions
+41
View File
@@ -589,6 +589,47 @@ When dynamic profile switching is enabled and a profile is activated via `gitea_
2. Call `gitea_whoami` with the target remote to prove and verify the fresh Gitea authenticated identity. 2. Call `gitea_whoami` with the target remote to prove and verify the fresh Gitea authenticated identity.
This guarantees the active profile operations align with the actual Gitea authenticated user credential. This guarantees the active profile operations align with the actual Gitea authenticated user credential.
### 4. Review-State Invalidation on Profile Switch (#690)
A cross-profile activation (e.g. reviewer → author → reviewer) is a session
boundary for formal review state. On any switch where the activated profile
differs from the previous one, `gitea_activate_profile` invalidates, in
memory **and** in durable session state (for both the old and new profile
identities):
- preflight identity/capability stamps (`gitea_whoami` / `gitea_resolve_task_capability` proof),
- review workflow-load proof (`gitea_load_review_workflow`),
- the review decision lock (including `final_review_decision_ready` markers),
- the reviewer PR session lease binding,
- live namespace-health assessments.
Before any formal verdict (`gitea_mark_final_review_decision` /
`gitea_submit_pr_review`) the full reviewer preflight must be re-established
under the new profile: `gitea_whoami`, `gitea_load_review_workflow`,
`gitea_resolve_task_capability(review_pr)`, live head re-pin, and lease
re-acquire/adopt. Switching back to the earlier profile does **not**
resurrect the prior run — durable state keyed by either profile identity is
cleared at switch time.
The supported pattern remains **separate session/namespace per role**
(dual-namespace, §2): file author-side follow-ups from an author session,
not by hopping profiles inside a formal review run. Runtime profile
switching is the operator-approved exception and always carries the
re-preflight cost above.
### 5. Namespace Provenance (#690)
The server cannot derive its own client-managed MCP namespace name, so a
launcher may declare it via the `GITEA_MCP_NAMESPACE` environment variable
(e.g. `gitea-reviewer`). `gitea_whoami`, `gitea_get_runtime_context`, and
`gitea_resolve_task_capability` report `namespace_provenance` — the
configured client namespace, the active execution profile, and, for tasks
with a required namespace (`review_pr` → `gitea-reviewer`, `merge_pr` →
`gitea-merger`), a mismatch verdict. A declared namespace that disagrees
with the requested task's required namespace **fails closed**. An
undeclared namespace is reported as `unknown` and is never treated as
proof either way.
## Gitea MCP Runtime Isolation and Worktree Safety ## Gitea MCP Runtime Isolation and Worktree Safety
To ensure high availability and prevent broken feature worktrees from disabling essential security/identity controls, the Gitea MCP server implements runtime isolation: To ensure high availability and prevent broken feature worktrees from disabling essential security/identity controls, the Gitea MCP server implements runtime isolation:
@@ -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.
+21
View File
@@ -86,3 +86,24 @@ When a namespace returns EOF, follow
When blocked, repair the IDE namespace and re-record a healthy When blocked, repair the IDE namespace and re-record a healthy
`client_namespace` assessment before retrying the mutation. `client_namespace` assessment before retrying the mutation.
## Namespace provenance (#690)
A server process cannot derive the name of the client-managed namespace it is
registered under, so the launcher may declare it with the
`GITEA_MCP_NAMESPACE` environment variable (e.g. `GITEA_MCP_NAMESPACE=gitea-reviewer`).
- `gitea_whoami`, `gitea_get_runtime_context`, and
`gitea_resolve_task_capability` report `namespace_provenance`: the declared
client namespace, the active execution profile, and — for tasks with a
required namespace (`review_pr`/`submit_review``gitea-reviewer`,
`merge_pr``gitea-merger`) — a `mismatch` verdict.
- A declared namespace that disagrees with the requested task's required
namespace **fails closed** (`allowed_in_current_session=false` with a STOP
guidance entry).
- An undeclared namespace is reported as `namespace_source="unknown"` and is
never treated as proof either way.
- A profile switch via `gitea_activate_profile` clears all recorded live
namespace-health assessments; re-probe through the client before further
review/merge mutations.
+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.
+119 -1
View File
@@ -839,6 +839,75 @@ def _invalidate_preflight_identity_state() -> None:
_clear_preflight_capability_state() _clear_preflight_capability_state()
# #690: session-boundary invalidation record for the most recent cross-profile
# activation. Surfaced in runtime diagnostics so a formal review run can prove
# its state was reset by a profile switch and must be fully re-established.
_PROFILE_SWITCH_INVALIDATION: dict | None = None
def _invalidate_review_state_on_profile_switch(
before_profile: str,
after_profile: str,
) -> dict:
"""Invalidate review-critical session state on a profile switch (#690).
Workflow-load proof, reviewer lease binding, the review decision lock,
live namespace health, and preflight identity/capability stamps recorded
under the prior profile are contaminated for the new role. Durable state
keyed by *either* profile identity is cleared so a reviewer author
reviewer hop cannot resurrect a stale review run: the full reviewer
preflight (gitea_whoami, gitea_load_review_workflow,
gitea_resolve_task_capability(review_pr), live head re-pin, lease
re-acquire/adopt) must be re-established before any formal verdict.
"""
global _PROFILE_SWITCH_INVALIDATION
invalidated: list[str] = []
_invalidate_preflight_identity_state()
invalidated.append("preflight_identity_capability")
review_workflow_load.clear_review_workflow_load()
invalidated.append("review_workflow_load")
_save_review_decision_lock(None)
invalidated.append("review_decision_lock")
reviewer_pr_lease.clear_session_lease()
invalidated.append("reviewer_session_lease")
if _LIVE_NAMESPACE_HEALTH:
_LIVE_NAMESPACE_HEALTH.clear()
invalidated.append("live_namespace_health")
# Durable records keyed by either profile identity must not survive the
# switch, or activating author → reviewer → author could revive a stale
# review run without re-preflight.
for identity in {before_profile, after_profile}:
if not identity:
continue
try:
mcp_session_state.clear_state(
kind=mcp_session_state.KIND_DECISION_LOCK,
profile_identity=identity,
)
mcp_session_state.clear_state(
kind=mcp_session_state.KIND_WORKFLOW_LOAD,
profile_identity=identity,
)
except Exception:
pass # best-effort durable cleanup; in-memory state already reset
invalidated.append("durable_profile_state")
_PROFILE_SWITCH_INVALIDATION = {
"from_profile": before_profile,
"to_profile": after_profile,
"invalidated": invalidated,
"invalidated_at": datetime.now(timezone.utc).isoformat(),
"re_preflight_required": True,
}
return dict(_PROFILE_SWITCH_INVALIDATION)
def record_preflight_check( def record_preflight_check(
type_name: str, type_name: str,
resolved_role: str | None = None, resolved_role: str | None = None,
@@ -17210,6 +17279,11 @@ def gitea_whoami(
"session_context_audit": session_ctx.mutation_context_audit_fields(), "session_context_audit": session_ctx.mutation_context_audit_fields(),
"identity_match": not id_match.get("block"), "identity_match": not id_match.get("block"),
"identity_match_reasons": id_match.get("reasons") or [], "identity_match_reasons": id_match.get("reasons") or [],
# #690 AC4: report launcher-declared client namespace alongside the
# active execution profile so drift is visible in diagnostics.
"namespace_provenance": mcp_namespace_health.namespace_provenance(
active_profile=profile["profile_name"]
),
} }
if id_match.get("block"): if id_match.get("block"):
_invalidate_preflight_identity_state() _invalidate_preflight_identity_state()
@@ -18108,6 +18182,11 @@ def gitea_get_runtime_context(
"shell_health": native_mcp_preference.shell_health_status(), "shell_health": native_mcp_preference.shell_health_status(),
"workflow_load_proof": review_workflow_load.workflow_load_status( "workflow_load_proof": review_workflow_load.workflow_load_status(
PROJECT_ROOT), PROJECT_ROOT),
# #690: namespace provenance + profile-switch invalidation evidence.
"namespace_provenance": mcp_namespace_health.namespace_provenance(
active_profile=profile["profile_name"]
),
"profile_switch_invalidation": _PROFILE_SWITCH_INVALIDATION,
} }
# #702: read-only visibility into the inherited GITEA_ACTIVE_WORKTREE # #702: read-only visibility into the inherited GITEA_ACTIVE_WORKTREE
@@ -18611,6 +18690,17 @@ def gitea_activate_profile(
source="gitea_activate_profile", source="gitea_activate_profile",
) )
# 4.7 #690: a profile switch is a session-boundary event for review state.
# Any workflow-load proof, reviewer lease, decision lock, namespace
# health, or preflight stamp recorded under the prior profile is
# contaminated for the new role and must be re-established under the new
# profile before any formal review verdict.
switch_invalidation = None
if before_profile != after_profile:
switch_invalidation = _invalidate_review_state_on_profile_switch(
before_profile, after_profile
)
# 5. Audit the switch if auditing is on # 5. Audit the switch if auditing is on
_audit( _audit(
"activate_profile", "activate_profile",
@@ -18621,11 +18711,12 @@ def gitea_activate_profile(
"before": before_profile, "before": before_profile,
"after": after_profile, "after": after_profile,
"session_context": session_ctx.mutation_context_audit_fields(), "session_context": session_ctx.mutation_context_audit_fields(),
"review_state_invalidated": bool(switch_invalidation),
}, },
username=after_identity, username=after_identity,
) )
return { result = {
"success": True, "success": True,
"message": f"Successfully activated profile '{profile_name}' (fresh identity verification complete).", "message": f"Successfully activated profile '{profile_name}' (fresh identity verification complete).",
"before_profile": before_profile, "before_profile": before_profile,
@@ -18635,6 +18726,18 @@ def gitea_activate_profile(
"session_context_audit": session_ctx.mutation_context_audit_fields(), "session_context_audit": session_ctx.mutation_context_audit_fields(),
"auto_profile_substitution": False, "auto_profile_substitution": False,
} }
if switch_invalidation is not None:
result["review_state_invalidation"] = switch_invalidation
result["re_preflight_required"] = True
result["exact_next_action"] = (
"Profile switch invalidated workflow-load proof, reviewer lease, "
"decision lock, and preflight stamps (#690). Before any formal "
"review verdict, re-run the full reviewer preflight: "
"gitea_whoami, gitea_load_review_workflow, "
"gitea_resolve_task_capability(review_pr), live head re-pin, and "
"lease re-acquire/adopt."
)
return result
@mcp.tool() @mcp.tool()
@@ -20879,12 +20982,22 @@ def gitea_resolve_task_capability(
f"{required_role} task '{task}' even if nearby permissions are " f"{required_role} task '{task}' even if nearby permissions are "
"present (fail closed)." "present (fail closed)."
) )
# #690 AC4: when the launcher declares a client namespace, a task with a
# required namespace must fail closed on mismatch (e.g. review_pr served
# from an author namespace).
ns_provenance = mcp_namespace_health.namespace_provenance(
task=task_key, active_profile=profile.get("profile_name")
)
ns_mismatch_reason = None
if ns_provenance.get("mismatch"):
ns_mismatch_reason = "; ".join(ns_provenance.get("reasons") or [])
cross_host_block = bool(remote_assess.get("block")) cross_host_block = bool(remote_assess.get("block"))
identity_block = bool(id_assess.get("block")) identity_block = bool(id_assess.get("block"))
drift_block = bool(ctx_assess.get("block")) drift_block = bool(ctx_assess.get("block"))
allowed_in_current_session = ( allowed_in_current_session = (
permission_allowed_in_current_session permission_allowed_in_current_session
and role_matches_current_session and role_matches_current_session
and not ns_provenance.get("mismatch")
and not cross_host_block and not cross_host_block
and not identity_block and not identity_block
and not drift_block and not drift_block
@@ -20935,6 +21048,8 @@ def gitea_resolve_task_capability(
) )
if role_mismatch_reason: if role_mismatch_reason:
deny_parts.append(role_mismatch_reason) deny_parts.append(role_mismatch_reason)
if ns_mismatch_reason:
deny_parts.append(ns_mismatch_reason)
if deny_parts: if deny_parts:
reason_msg = "; ".join(deny_parts) reason_msg = "; ".join(deny_parts)
elif configured and switching: elif configured and switching:
@@ -21012,6 +21127,8 @@ def gitea_resolve_task_capability(
task_role_guidance = [] task_role_guidance = []
if role_mismatch_reason: if role_mismatch_reason:
task_role_guidance.append(f"STOP: {role_mismatch_reason}") task_role_guidance.append(f"STOP: {role_mismatch_reason}")
if ns_mismatch_reason:
task_role_guidance.append(f"STOP: {ns_mismatch_reason}")
if required_role == "reviewer": if required_role == "reviewer":
if allowed_in_current_session: if allowed_in_current_session:
task_role_guidance.append( task_role_guidance.append(
@@ -21078,6 +21195,7 @@ def gitea_resolve_task_capability(
"session_context_audit": session_ctx.mutation_context_audit_fields(), "session_context_audit": session_ctx.mutation_context_audit_fields(),
"profile_remote_compatible": not cross_host_block, "profile_remote_compatible": not cross_host_block,
"identity_match": not identity_block, "identity_match": not identity_block,
"namespace_provenance": ns_provenance,
"auto_profile_substitution": False, "auto_profile_substitution": False,
} }
# #685: report typed reconnect blocker without mutating config or exiting. # #685: report typed reconnect blocker without mutating config or exiting.
+49
View File
@@ -16,6 +16,7 @@ Probe sources
from __future__ import annotations from __future__ import annotations
import os
from typing import Any from typing import Any
@@ -57,8 +58,56 @@ SAFE_ENV_KEYS = (
"GITEA_SERVICE", "GITEA_SERVICE",
"GITEA_EXECUTION_ROLE", "GITEA_EXECUTION_ROLE",
"GITEA_MCP_CONFIG", "GITEA_MCP_CONFIG",
"GITEA_MCP_NAMESPACE",
) )
# Optional launcher-provided env declaring the client-managed MCP namespace
# this process is registered under (e.g. ``gitea-reviewer``). The server
# cannot derive its own IDE namespace name, so the launcher declares it; when
# declared, reviewers/mergers can fail closed on a namespace/task mismatch
# (#690 AC4). Absence means "unknown" — reported, never guessed.
NAMESPACE_ENV = "GITEA_MCP_NAMESPACE"
def configured_client_namespace(env: dict[str, str] | None = None) -> str | None:
"""Return the launcher-declared client namespace, or None when unknown."""
source = os.environ if env is None else env
value = (source.get(NAMESPACE_ENV) or "").strip()
return value or None
def namespace_provenance(
task: str | None = None,
*,
active_profile: str | None = None,
env: dict[str, str] | None = None,
) -> dict[str, Any]:
"""Report configured client namespace vs active execution profile (#690).
When *task* carries a required namespace (``TASK_REQUIRED_NAMESPACES``)
and the launcher declared a different one, ``mismatch`` is True and the
caller must fail closed for that task. An undeclared namespace is
reported as unknown — never treated as proof either way.
"""
configured = configured_client_namespace(env)
required = TASK_REQUIRED_NAMESPACES.get(task or "")
mismatch = bool(configured and required and configured != required)
reasons: list[str] = []
if mismatch:
reasons.append(
f"configured client namespace '{configured}' does not match "
f"required namespace '{required}' for task '{task}' (fail closed)"
)
return {
"configured_namespace": configured,
"namespace_source": NAMESPACE_ENV if configured else "unknown",
"active_profile": active_profile,
"requested_task": task,
"required_namespace": required,
"mismatch": mismatch,
"reasons": reasons,
}
def _as_list(value: Any) -> list[str] | None: def _as_list(value: Any) -> list[str] | None:
if value is None: if value is None:
+2
View File
@@ -41,6 +41,7 @@ def _reset_mutation_authority(monkeypatch):
"GITEA_REVIEWER_WORKTREE", "GITEA_REVIEWER_WORKTREE",
"GITEA_MERGER_WORKTREE", "GITEA_MERGER_WORKTREE",
"GITEA_RECONCILER_WORKTREE", "GITEA_RECONCILER_WORKTREE",
"GITEA_MCP_NAMESPACE",
]: ]:
monkeypatch.delenv(env_key, raising=False) monkeypatch.delenv(env_key, raising=False)
@@ -115,6 +116,7 @@ def _reset_mutation_authority(monkeypatch):
monkeypatch.setattr(mcp_server, "_ACTOR_IDENTITY_CACHE", {}) monkeypatch.setattr(mcp_server, "_ACTOR_IDENTITY_CACHE", {})
monkeypatch.setattr(mcp_server, "_REVIEW_DECISION_LOCK", None) monkeypatch.setattr(mcp_server, "_REVIEW_DECISION_LOCK", None)
monkeypatch.setattr(mcp_server, "_LIVE_NAMESPACE_HEALTH", {}) monkeypatch.setattr(mcp_server, "_LIVE_NAMESPACE_HEALTH", {})
monkeypatch.setattr(mcp_server, "_PROFILE_SWITCH_INVALIDATION", None)
monkeypatch.setattr(mcp_server, "_preflight_whoami_called", False) monkeypatch.setattr(mcp_server, "_preflight_whoami_called", False)
monkeypatch.setattr(mcp_server, "_preflight_capability_called", False) monkeypatch.setattr(mcp_server, "_preflight_capability_called", False)
monkeypatch.setattr(mcp_server, "_preflight_resolved_role", None) monkeypatch.setattr(mcp_server, "_preflight_resolved_role", None)
@@ -0,0 +1,274 @@
"""Regression coverage for #690: cross-role profile activation invalidation.
A mid-run profile switch (e.g. reviewer → author → reviewer) must invalidate
workflow-load proof, reviewer lease binding, review decision lock, live
namespace health, and preflight identity/capability stamps, and must require
a full reviewer preflight before any formal verdict. Namespace provenance
must be reported and fail closed on task/namespace mismatch.
"""
import json
import os
import sys
import tempfile
import unittest
from unittest.mock import patch
sys.path.insert(0, str(__import__("pathlib").Path(__file__).resolve().parent.parent))
import gitea_config
import mcp_namespace_health
import mcp_server
import mcp_session_state
import review_workflow_load
import reviewer_pr_lease
from tests.test_runtime_clarity import CONFIG_SWITCHING_ENABLED
class TestProfileSwitchReviewGuard(unittest.TestCase):
def setUp(self):
self._remotes_patch = patch.dict(mcp_server.REMOTES, {
"dadeschools": {"host": "gitea.example.com", "org": "Example-Org", "repo": "Example-Repo"},
"prgs": {"host": "gitea.example.com", "org": "Example-Org", "repo": "Example-Repo"},
})
self._remotes_patch.start()
mcp_server._IDENTITY_CACHE.clear()
gitea_config._active_profile_override = None
self._dir = tempfile.TemporaryDirectory()
self.config_path = os.path.join(self._dir.name, "profiles.json")
with open(self.config_path, "w", encoding="utf-8") as fh:
fh.write(json.dumps(CONFIG_SWITCHING_ENABLED))
def tearDown(self):
self._remotes_patch.stop()
mcp_server._IDENTITY_CACHE.clear()
gitea_config._active_profile_override = None
self._dir.cleanup()
def _env(self, profile="reviewer-profile"):
return {
"GITEA_MCP_CONFIG": self.config_path,
"GITEA_MCP_PROFILE": profile,
"GITEA_TOKEN_AUTHOR": "author-pass",
"GITEA_TOKEN_REVIEWER": "reviewer-pass",
"GITEA_TOKEN_MERGER": "merger-pass",
}
def _seed_contaminated_review_state(self):
"""Simulate an in-flight reviewer run under reviewer-profile."""
mcp_server._preflight_whoami_called = True
mcp_server._preflight_capability_called = True
mcp_server._preflight_resolved_role = "reviewer"
mcp_server._preflight_resolved_task = "review_pr"
review_workflow_load._REVIEW_WORKFLOW_LOAD = {"loaded": True}
mcp_server._REVIEW_DECISION_LOCK = {
"session_profile": "reviewer-profile",
"final_review_decision_ready": True,
"ready_pr_number": 688,
}
reviewer_pr_lease.record_session_lease(
{"session_id": "lease-session-1", "pr_number": 688}
)
mcp_server._LIVE_NAMESPACE_HEALTH["gitea-reviewer"] = {
"namespace": "gitea-reviewer",
"healthy": True,
"ide_namespace_proven": True,
}
# Durable records keyed by the reviewer identity must also be cleared.
mcp_session_state.save_state(
kind=mcp_session_state.KIND_WORKFLOW_LOAD,
payload={"loaded": True},
profile_identity="reviewer-profile",
)
mcp_session_state.save_state(
kind=mcp_session_state.KIND_DECISION_LOCK,
payload={"final_review_decision_ready": True, "ready_pr_number": 688},
profile_identity="reviewer-profile",
)
def _activate(self, target, logins):
with patch.object(
mcp_server, "get_auth_header", side_effect=[f"token p" for _ in logins]
), patch.object(
mcp_server, "api_request", side_effect=[{"login": l} for l in logins]
), patch.object(
mcp_server,
"_workspace_repository_slug",
return_value="Example-Org/Example-Repo",
), patch.object(
mcp_server, "_canonical_repository_slug", return_value=(None, [])
):
return mcp_server.gitea_activate_profile(profile_name=target)
# -----------------------------------------------------------------
# AC1/AC2/AC3: switch invalidates review state; re-preflight required
# -----------------------------------------------------------------
def test_switch_invalidates_review_state_and_blocks_verdict(self):
with patch.dict(os.environ, self._env("reviewer-profile"), clear=True):
self._seed_contaminated_review_state()
res = self._activate("author-profile", ["reviewer-user", "author-user"])
self.assertTrue(res["success"])
self.assertTrue(res["re_preflight_required"])
inv = res["review_state_invalidation"]
self.assertEqual(inv["from_profile"], "reviewer-profile")
self.assertEqual(inv["to_profile"], "author-profile")
for item in (
"preflight_identity_capability",
"review_workflow_load",
"review_decision_lock",
"reviewer_session_lease",
"live_namespace_health",
):
self.assertIn(item, inv["invalidated"])
# In-memory state cleared.
self.assertFalse(mcp_server._preflight_whoami_called)
self.assertFalse(mcp_server._preflight_capability_called)
self.assertIsNone(mcp_server._preflight_resolved_task)
self.assertIsNone(review_workflow_load._REVIEW_WORKFLOW_LOAD)
self.assertIsNone(mcp_server._REVIEW_DECISION_LOCK)
self.assertIsNone(reviewer_pr_lease.get_session_lease())
self.assertEqual(mcp_server._LIVE_NAMESPACE_HEALTH, {})
self.assertIsNotNone(mcp_server._PROFILE_SWITCH_INVALIDATION)
# Durable records keyed by the reviewer identity are gone.
self.assertIsNone(
mcp_session_state.load_state(
kind=mcp_session_state.KIND_WORKFLOW_LOAD,
profile_identity="reviewer-profile",
)
)
self.assertIsNone(
mcp_session_state.load_state(
kind=mcp_session_state.KIND_DECISION_LOCK,
profile_identity="reviewer-profile",
)
)
# A formal verdict without re-preflight fails closed.
reasons = mcp_server.check_review_decision_gate(
688, "APPROVE", final_review_decision_ready=True
)
self.assertTrue(reasons)
def test_switch_back_cannot_resurrect_stale_review_run(self):
with patch.dict(os.environ, self._env("reviewer-profile"), clear=True):
self._seed_contaminated_review_state()
self._activate("author-profile", ["reviewer-user", "author-user"])
res = self._activate("reviewer-profile", ["author-user", "reviewer-user"])
self.assertTrue(res["success"])
# The pre-switch review run must not reappear.
self.assertIsNone(mcp_server._REVIEW_DECISION_LOCK)
self.assertIsNone(review_workflow_load._REVIEW_WORKFLOW_LOAD)
self.assertIsNone(reviewer_pr_lease.get_session_lease())
status = review_workflow_load.workflow_load_status()
self.assertFalse(status["workflow_load_valid"])
reasons = mcp_server.check_review_decision_gate(
688, "APPROVE", final_review_decision_ready=True
)
self.assertTrue(reasons)
def test_same_profile_reactivation_keeps_state(self):
with patch.dict(os.environ, self._env("reviewer-profile"), clear=True):
self._seed_contaminated_review_state()
res = self._activate("reviewer-profile", ["reviewer-user", "reviewer-user"])
self.assertTrue(res["success"], res)
self.assertNotIn("review_state_invalidation", res)
self.assertIsNotNone(mcp_server._REVIEW_DECISION_LOCK)
self.assertTrue(mcp_server._preflight_whoami_called)
def test_clean_repreflight_after_switch_allows_gate(self):
with patch.dict(os.environ, self._env("reviewer-profile"), clear=True):
self._seed_contaminated_review_state()
self._activate("author-profile", ["reviewer-user", "author-user"])
self._activate("reviewer-profile", ["author-user", "reviewer-user"])
# Re-establish the full reviewer preflight under the new profile.
mcp_server.record_preflight_check("whoami")
mcp_server.record_preflight_check(
"capability", resolved_role="reviewer", resolved_task="review_pr"
)
mcp_server.init_review_decision_lock("dadeschools", "review_pr")
lock = mcp_server._load_review_decision_lock()
self.assertIsNotNone(lock)
lock.update(
{
"final_review_decision_ready": True,
"ready_pr_number": 688,
"ready_action": "APPROVE",
"ready_remote": "dadeschools",
"ready_org": "Example-Org",
"ready_repo": "Example-Repo",
}
)
mcp_server._save_review_decision_lock(lock)
with patch.object(
mcp_server, "_review_workflow_load_gate_reasons", return_value=[]
):
reasons = mcp_server.check_review_decision_gate(
688,
"APPROVE",
final_review_decision_ready=True,
remote="dadeschools",
)
self.assertEqual(reasons, [])
# -----------------------------------------------------------------
# AC4: namespace provenance reporting + fail-closed mismatch
# -----------------------------------------------------------------
def test_namespace_provenance_mismatch_detection(self):
prov = mcp_namespace_health.namespace_provenance(
task="review_pr",
active_profile="reviewer-profile",
env={"GITEA_MCP_NAMESPACE": "gitea-author"},
)
self.assertTrue(prov["mismatch"])
self.assertEqual(prov["required_namespace"], "gitea-reviewer")
prov_ok = mcp_namespace_health.namespace_provenance(
task="review_pr",
active_profile="reviewer-profile",
env={"GITEA_MCP_NAMESPACE": "gitea-reviewer"},
)
self.assertFalse(prov_ok["mismatch"])
prov_unknown = mcp_namespace_health.namespace_provenance(
task="review_pr", active_profile="reviewer-profile", env={}
)
self.assertIsNone(prov_unknown["configured_namespace"])
self.assertFalse(prov_unknown["mismatch"])
self.assertEqual(prov_unknown["namespace_source"], "unknown")
@patch("mcp_server.api_request", return_value={"login": "reviewer-user"})
@patch("mcp_server.get_auth_header", return_value="token reviewer-pass")
def test_whoami_reports_namespace_provenance(self, _auth, _api):
env = self._env("reviewer-profile")
env["GITEA_MCP_NAMESPACE"] = "gitea-reviewer"
with patch.dict(os.environ, env, clear=True):
res = mcp_server.gitea_whoami(remote="dadeschools")
prov = res["namespace_provenance"]
self.assertEqual(prov["configured_namespace"], "gitea-reviewer")
self.assertEqual(prov["active_profile"], "reviewer-profile")
self.assertFalse(prov["mismatch"])
@patch("mcp_server.api_request", return_value={"login": "reviewer-user"})
@patch("mcp_server.get_auth_header", return_value="token reviewer-pass")
def test_resolve_fails_closed_on_namespace_mismatch(self, _auth, _api):
env = self._env("reviewer-profile")
env["GITEA_MCP_NAMESPACE"] = "gitea-author"
with patch.dict(os.environ, env, clear=True):
res = mcp_server.gitea_resolve_task_capability(
task="review_pr", kwargs="{}", remote="dadeschools"
)
self.assertFalse(res["allowed_in_current_session"])
self.assertTrue(res["namespace_provenance"]["mismatch"])
self.assertTrue(
any("namespace" in g for g in res["task_role_guidance"])
)
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()
+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, snapshot_to_dict as system_health_to_dict,
) )
from webui.system_health_views import render_system_health_page 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 import request_service
from webui.request_views import render_requests_page 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]: def _default_request_scope() -> dict[str, str]:
"""Resolve remote/org/repo from the project registry for request forms. """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("/api/queue", api_queue, methods=["GET"]),
Route("/traffic", traffic, methods=["GET"]), Route("/traffic", traffic, methods=["GET"]),
Route("/api/traffic", api_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", projects, methods=["GET"]),
Route("/projects/{project_id}", project_detail, methods=["GET"]), Route("/projects/{project_id}", project_detail, methods=["GET"]),
Route("/api/projects", api_projects, methods=["GET"]), Route("/api/projects", api_projects, methods=["GET"]),
+1
View File
@@ -46,6 +46,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
NavItem("/queue", "Queue"), NavItem("/queue", "Queue"),
NavItem("/leases", "Leases"), NavItem("/leases", "Leases"),
NavItem("/actions", "Actions"), NavItem("/actions", "Actions"),
NavItem("/notifications", "Notifications"),
NavItem("/requests", "Requests"), NavItem("/requests", "Requests"),
)), )),
NavGroup("Runtime/Sessions", ( 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()