Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b993ad1c64 | ||
|
|
d0006e9f71 | ||
|
|
6da68fffb8 | ||
|
|
53ce1b1a5e | ||
|
|
433f66add8 |
+98
-24
@@ -23,6 +23,7 @@ import json
|
||||
import os
|
||||
import uuid
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Mapping, Sequence
|
||||
|
||||
from control_plane_db import (
|
||||
@@ -738,6 +739,46 @@ def normalize_exclude_issue_numbers(
|
||||
return sorted(out)
|
||||
|
||||
|
||||
def _claim_expires_at(claim: Any) -> datetime | None:
|
||||
"""Parse a claim's ``expires_at``, or ``None`` when it is absent/malformed."""
|
||||
if not isinstance(claim, Mapping):
|
||||
return None
|
||||
text = str(claim.get("expires_at") or "").strip()
|
||||
if not text:
|
||||
return None
|
||||
if text.endswith("Z"):
|
||||
text = text[:-1] + "+00:00"
|
||||
try:
|
||||
parsed = datetime.fromisoformat(text)
|
||||
except ValueError:
|
||||
return None
|
||||
if parsed.tzinfo is None:
|
||||
parsed = parsed.replace(tzinfo=timezone.utc)
|
||||
return parsed.astimezone(timezone.utc)
|
||||
|
||||
|
||||
def _drop_expired_claims(
|
||||
claims: Mapping[tuple[str, int], dict[str, Any]],
|
||||
*,
|
||||
now: datetime | None = None,
|
||||
) -> dict[tuple[str, int], dict[str, Any]]:
|
||||
"""Claims minus those whose lease has already expired (#643).
|
||||
|
||||
The read-only mirror of ``expire_stale_leases``: the sweep marks such rows
|
||||
``expired`` so they stop being returned as claims, and this reaches the same
|
||||
view without writing. A claim with no parseable ``expires_at`` is **kept** —
|
||||
an unreadable expiry is not evidence that work is free.
|
||||
"""
|
||||
moment = now or datetime.now(timezone.utc)
|
||||
kept: dict[tuple[str, int], dict[str, Any]] = {}
|
||||
for key, claim in (claims or {}).items():
|
||||
expires_at = _claim_expires_at(claim)
|
||||
if expires_at is not None and expires_at <= moment:
|
||||
continue
|
||||
kept[key] = claim
|
||||
return kept
|
||||
|
||||
|
||||
def candidate_set_fingerprint(
|
||||
candidates: Sequence[WorkCandidate],
|
||||
*,
|
||||
@@ -826,12 +867,22 @@ def allocate_next_work(
|
||||
exclude_issue_numbers: Sequence[int] | None = None,
|
||||
expected_candidate_set_fingerprint: str | None = None,
|
||||
allocation_mode: str | None = None,
|
||||
side_effect_free: bool = False,
|
||||
) -> dict[str, Any]:
|
||||
"""Select and optionally reserve the next work unit via control-plane DB.
|
||||
|
||||
*apply=False* (default): dry-run selection only — no lease/assignment.
|
||||
*apply=True*: atomic ``assign_and_lease`` for the selected candidate.
|
||||
|
||||
*side_effect_free* (#643): a dry run that writes **nothing** to the
|
||||
control-plane DB. A plain ``apply=False`` still registered a session row and
|
||||
swept stale leases globally, so a caller advertising a read-only preview was
|
||||
mutating on every call. Under this flag both writes are suppressed and stale
|
||||
leases are instead filtered out of the claim map in memory, which yields the
|
||||
same selection the sweep would have produced without persisting anything.
|
||||
Incompatible with *apply* — the combination fails closed rather than
|
||||
silently reserving.
|
||||
|
||||
*allocation_mode* (#840): ``cross_role`` (default for controller) inspects
|
||||
the complete queue and returns one authoritative selection naming the
|
||||
required downstream role/profile/action. ``role_scoped`` keeps prior
|
||||
@@ -885,40 +936,57 @@ def allocate_next_work(
|
||||
"allocation_mode": (allocation_mode or "").strip() or None,
|
||||
}
|
||||
|
||||
session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}"
|
||||
try:
|
||||
db.upsert_session(
|
||||
session_id=session_id,
|
||||
role=role_norm,
|
||||
profile=profile_name,
|
||||
pid=os.getpid(),
|
||||
controller_instance_id=controller_instance_id,
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 — surface structured
|
||||
# A side-effect-free run may never reserve: reserving is a write, and the
|
||||
# flag is the caller's assertion that this call writes nothing (#643).
|
||||
if side_effect_free and apply:
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"apply": True,
|
||||
"reasons": [
|
||||
f"failed to register session in control-plane DB: {exc} "
|
||||
"(fail closed, #613)"
|
||||
"side_effect_free is incompatible with apply=True; an "
|
||||
"assignment is a write (fail closed, #643)"
|
||||
],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
|
||||
# Expire stale leases globally before selection.
|
||||
try:
|
||||
db.expire_stale_leases()
|
||||
except Exception as exc: # noqa: BLE001
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"reasons": [f"lease expiry failed: {exc} (fail closed)"],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}"
|
||||
if not side_effect_free:
|
||||
try:
|
||||
db.upsert_session(
|
||||
session_id=session_id,
|
||||
role=role_norm,
|
||||
profile=profile_name,
|
||||
pid=os.getpid(),
|
||||
controller_instance_id=controller_instance_id,
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 — surface structured
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"reasons": [
|
||||
f"failed to register session in control-plane DB: {exc} "
|
||||
"(fail closed, #613)"
|
||||
],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
|
||||
# Expire stale leases globally before selection.
|
||||
try:
|
||||
db.expire_stale_leases()
|
||||
except Exception as exc: # noqa: BLE001
|
||||
return {
|
||||
"success": False,
|
||||
"outcome": OUTCOME_NO_SAFE,
|
||||
"reasons": [f"lease expiry failed: {exc} (fail closed)"],
|
||||
"skipped": [],
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
|
||||
terminal = None
|
||||
try:
|
||||
@@ -953,6 +1021,12 @@ def allocate_next_work(
|
||||
"assignment": None,
|
||||
"substrate": "control_plane_db",
|
||||
}
|
||||
if side_effect_free:
|
||||
# ``list_active_claims`` filters on status alone, so without the
|
||||
# global sweep an already-expired lease would still read as a live
|
||||
# claim and the preview would report work as taken that is free.
|
||||
# Drop those in memory: same view the sweep produces, no write.
|
||||
claims = _drop_expired_claims(claims)
|
||||
|
||||
try:
|
||||
exclude_nums = normalize_exclude_issue_numbers(exclude_issue_numbers)
|
||||
|
||||
@@ -94,6 +94,7 @@ already define, and a regression test asserts each mapping matches.
|
||||
| `record_analytics_usage` | operator | gated_write | `runtime.record_analytics_usage` | Yes | No | No | 2 |
|
||||
| `system.reload_namespace` | controller | privileged | `runtime.reload_namespace` | Yes | No | No | 2 |
|
||||
| `system.restart_namespace` | admin | destructive | `runtime.restart_namespace` | Yes | **Yes** | **Yes** | 2 |
|
||||
| `initiate_workflow` | operator | gated_write | `gitea.read` | Yes | No | No | 2 |
|
||||
|
||||
**Dual control** means the acting principal may not be the sole authority: a
|
||||
second distinct principal must confirm. **Break-glass** means the action is
|
||||
@@ -112,6 +113,12 @@ by the console — both hand off to a host supervisor, and neither exposes a raw
|
||||
process kill. See
|
||||
[`sanctioned-restart-controls.md`](sanctioned-restart-controls.md) (#642).
|
||||
|
||||
`initiate_workflow` (#643) is operator-class because its outcome is a *claim*,
|
||||
not a Gitea verdict. Requesting reviewer or merger work reserves that work
|
||||
through the allocator; it does not grant the right to approve or merge, which
|
||||
stays with the MCP role profile and its own capability gates. See
|
||||
[`webui-requests.md`](webui-requests.md).
|
||||
|
||||
### Authorization decision
|
||||
|
||||
`authorize(action_id, principal, for_execution=False)` returns a decision
|
||||
@@ -126,9 +133,24 @@ record and **denies by default**. The deny reasons are closed and enumerated:
|
||||
| `phase_not_active` | Execution requested for an action whose phase is not open. |
|
||||
| `allowed_preview_only` | Authorized — preview only, execution still disabled. |
|
||||
|
||||
There is no implicit allow branch. Even the allow result reports
|
||||
`execution_enabled: false` while the console is in Phase 1, so no caller can
|
||||
read an allow as permission to mutate.
|
||||
There is no implicit allow branch.
|
||||
|
||||
`execution_enabled` on the decision reports whether the action has a live
|
||||
execution path at all, and is computed by `execution_wired(action)`. There are
|
||||
exactly two ways to be wired:
|
||||
|
||||
1. the action's `phase` is at or below `ACTIVE_PHASE`; or
|
||||
2. the action declares an `execution_env_flag` **and** that variable is set.
|
||||
|
||||
Every action that declares no flag therefore reports `execution_enabled: false`
|
||||
while the console is in Phase 1, so no caller can read an allow as permission
|
||||
to mutate. The per-action flag exists because raising `ACTIVE_PHASE` would
|
||||
enable execution for every action of that phase at once, including ones whose
|
||||
execution path is not implemented. One implemented action goes live on its own
|
||||
flag instead of dragging its unimplemented phase-mates with it.
|
||||
|
||||
`initiate_workflow` is the only action that currently declares a flag
|
||||
(`WEBUI_REQUESTS_EXECUTION`), and it stays denied until an operator sets it.
|
||||
|
||||
## Secret redaction
|
||||
|
||||
@@ -235,13 +257,22 @@ second one. The integration points are already wired and observable:
|
||||
instead of adding a parallel check.
|
||||
- **`GET /api/console/security-model`** publishes the RBAC matrix, redaction
|
||||
policy, and audit policy as JSON for operators and tests.
|
||||
- **`POST /api/v1/requests/preview` and `.../apply`** (#643) are the first
|
||||
actions to use this model for a real execution path. Preview always returns a
|
||||
decision and an audited `previewed` record; apply requires `confirm=true`,
|
||||
emits `succeeded` or `denied`, and reserves work only through the allocator.
|
||||
See [`webui-requests.md`](webui-requests.md).
|
||||
|
||||
To open Phase 2, a child issue must: raise `ACTIVE_PHASE`, implement the
|
||||
confirmation and dual-control flow the matrix already declares, emit a
|
||||
`succeeded` or `failed` record alongside the `gitea_audit` mutation record, and
|
||||
keep `viewer` unable to reach any of it. Turning on execution without the
|
||||
confirmation flow contradicts a declared requirement and is a review failure,
|
||||
not a shortcut.
|
||||
A Phase 2 action must: use `execution_wired` rather than a private enable flag,
|
||||
implement the confirmation and dual-control flow the matrix already declares,
|
||||
emit a `succeeded` or `failed` record alongside the `gitea_audit` mutation
|
||||
record, and keep `viewer` unable to reach any of it. Turning on execution
|
||||
without the confirmation flow contradicts a declared requirement and is a
|
||||
review failure, not a shortcut.
|
||||
|
||||
Raising `ACTIVE_PHASE` remains the way to open a whole phase at once, and is
|
||||
deliberately *not* what #643 did: an action-scoped opt-in cannot enable an
|
||||
action whose execution path nobody wrote.
|
||||
|
||||
## Local-dev mode
|
||||
|
||||
@@ -294,6 +325,7 @@ Until Phase 2 wires it, probe protection rests on network placement alone, as
|
||||
| `WEBUI_ROLE_MAP` | unset | JSON subject → role map |
|
||||
| `WEBUI_REQUIRE_PROBE_AUTH` | unset | Require auth for non-public probes |
|
||||
| `WEBUI_CONSOLE_AUDIT_LOG` | unset | Append-only audit sink path |
|
||||
| `WEBUI_REQUESTS_EXECUTION` | unset | Opt in to `initiate_workflow` execution (#643) |
|
||||
|
||||
All are read server-side only. None is ever rendered into a page or returned by
|
||||
an API.
|
||||
|
||||
@@ -1,81 +0,0 @@
|
||||
# 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.
|
||||
@@ -0,0 +1,160 @@
|
||||
# Web console requests: intent preview and workflow initiation (#643)
|
||||
|
||||
**Phase 2. Preview is always live and always read-only. Initiation is wired but
|
||||
denied until an operator opts in.**
|
||||
|
||||
Before this surface, starting role work meant pasting a prompt into a terminal
|
||||
and trusting the operator to have checked the allocator first. Nothing enforced
|
||||
that check, so two sessions could reach for the same issue and each believe it
|
||||
was theirs. This page replaces the paste with a *request*: a desired role, an
|
||||
issue or PR, and a stated intent, answered by an authorization decision and —
|
||||
on confirmation — an exclusive assignment from the allocator.
|
||||
|
||||
| Concern | Module |
|
||||
|---------|--------|
|
||||
| Request model, preview, initiation | `webui/request_service.py` |
|
||||
| Form and preview rendering | `webui/request_views.py` |
|
||||
| Authorization | `webui/console_authz.py` (`initiate_workflow`) |
|
||||
| Audit | `webui/console_audit.py` |
|
||||
| Ownership substrate | `allocator_service.py` + `control_plane_db.py` |
|
||||
|
||||
## Surfaces
|
||||
|
||||
| Path | Method | Purpose |
|
||||
|------|--------|---------|
|
||||
| `/requests` | GET | Request form |
|
||||
| `/requests` | POST | Render an intent preview. **Never assigns.** |
|
||||
| `/api/v1/requests/preview` | POST | Intent preview as JSON |
|
||||
| `/api/v1/requests/apply` | POST | Initiate — confirmed, audited, allocator-owned |
|
||||
|
||||
The HTML form has no initiate button on purpose. Initiating requires a
|
||||
confirmed POST to `/api/v1/requests/apply`, so a stray form submission cannot
|
||||
reserve work as a side effect.
|
||||
|
||||
## The request
|
||||
|
||||
```json
|
||||
{
|
||||
"desired_role": "author",
|
||||
"work_kind": "issue",
|
||||
"work_number": 643,
|
||||
"intent_summary": "implement request preview and initiation",
|
||||
"remote": "prgs",
|
||||
"org": "Scaled-Tech-Consulting",
|
||||
"repo": "Gitea-Tools",
|
||||
"expected_head_sha": null
|
||||
}
|
||||
```
|
||||
|
||||
`desired_role` is one of `author`, `reviewer`, `merger`, `reconciler`,
|
||||
`controller`. `work_kind` is `issue` or `pr`. `remote`/`org`/`repo` default to
|
||||
the first project in the registry when omitted; when neither the request nor
|
||||
the registry resolves them, the request is rejected rather than pointed at some
|
||||
other repository. `intent_summary` is required — it is what the audit record
|
||||
states as the reason — and is truncated to 500 characters.
|
||||
|
||||
Parsing rejects rather than corrects. An unknown role, an unknown work kind, a
|
||||
non-positive number, or a missing intent each return `400` with a `reason_code`
|
||||
and the offending `field`.
|
||||
|
||||
## Preview
|
||||
|
||||
Five checks, each with its own verdict, reason code, and detail:
|
||||
|
||||
| Check | Passes when |
|
||||
|-------|-------------|
|
||||
| `authorization` | The console principal holds `operator` or above |
|
||||
| `capability` | The desired role maps to a declared profile and MCP namespace |
|
||||
| `lease_availability` | No active claim holds the work unit |
|
||||
| `next_safe_action` | The allocator would independently select this exact work unit |
|
||||
| `head_pin` | PR work resolves to a head SHA, and a supplied SHA still matches |
|
||||
|
||||
A preview also returns the role's `allowed_actions` and `prohibited_actions`
|
||||
(from `allocator_service.ROLE_ACTIONS`), the `required_profile` and
|
||||
`required_namespace` the work must run under, and a `correlation_id` that ties
|
||||
the preview to its audit record and to any assignment that follows.
|
||||
|
||||
Preview is read-only in the strict sense: it calls the allocator with
|
||||
`apply=false` and writes nothing but an audit line. An unauthorized principal
|
||||
never reaches the allocator or the control-plane DB at all, so a denial cannot
|
||||
be used to enumerate the queue.
|
||||
|
||||
## Initiation
|
||||
|
||||
`POST /api/v1/requests/apply` refuses in this order, and every refusal returns
|
||||
before any assignment is attempted:
|
||||
|
||||
| Condition | Outcome | Status |
|
||||
|-----------|---------|--------|
|
||||
| Unparseable request | `invalid_request` | 400 |
|
||||
| Not authorized, or execution not wired | `denied` | 403 |
|
||||
| `confirm` not set | `denied` / `confirmation_required` | 409 |
|
||||
| Work unit already claimed | `blocked` / `duplicate_assignment` | 409 |
|
||||
| Allocator would select other work | `wait` / `not_next_safe_work` | 409 |
|
||||
| Allocator declines on apply | `blocked` or `wait` | 409 |
|
||||
| Evidence unavailable | `wait` / `evidence_unavailable` | 503 |
|
||||
| Assigned | `assigned_work` | 201 |
|
||||
|
||||
A success returns the assignment plus a `handoff` block naming the profile, the
|
||||
namespace, and the actions that stay forbidden — enough for the operator to
|
||||
continue in the right MCP namespace without guessing.
|
||||
|
||||
### Why apply runs the allocator twice
|
||||
|
||||
The allocator is the only source of exclusive ownership (#600 / #613), and it
|
||||
selects work; it does not take orders. So `apply` runs a dry-run first and
|
||||
proceeds only when the allocator would independently pick the requested work
|
||||
unit. If it would not, the request reports `wait` and mutates nothing.
|
||||
|
||||
A request is therefore a *confirmation* of the allocator's decision, never an
|
||||
override of it. The apply call carries the dry-run's
|
||||
`candidate_set_fingerprint` as a CAS pin (#776), so a queue that changed
|
||||
between the two calls fails closed rather than assigning against a stale view.
|
||||
The result is checked again on the way out: an assignment naming a different
|
||||
work unit is not read as success.
|
||||
|
||||
### Fail-closed defaults
|
||||
|
||||
- An unreadable control-plane DB denies. It is never treated as "nothing holds
|
||||
this work unit".
|
||||
- An incomplete queue inventory denies (#758). Ranking a partial candidate set
|
||||
can select the wrong work.
|
||||
- An allocator that raises denies.
|
||||
- PR work with no resolvable head SHA denies; a supplied SHA that no longer
|
||||
matches denies with `head_moved`.
|
||||
|
||||
## Enabling initiation
|
||||
|
||||
Execution is wired off. Set `WEBUI_REQUESTS_EXECUTION=1` to enable it for the
|
||||
`initiate_workflow` action only — see
|
||||
[`webui-authz-audit.md`](webui-authz-audit.md) for why this is an
|
||||
action-scoped flag rather than a phase bump. With the variable unset, `apply`
|
||||
returns `403` with `reason_code: unauthorized` no matter who asks.
|
||||
|
||||
Enabling execution does **not** enable approvals or merges. Those are phase 3
|
||||
console actions and remain forbidden in every path here; the console reserves
|
||||
work and hands off, and the MCP role profile enforces what that role may then
|
||||
do.
|
||||
|
||||
## Audit
|
||||
|
||||
Every preview and every apply emits a console audit record (schema in
|
||||
[`webui-authz-audit.md`](webui-authz-audit.md)):
|
||||
|
||||
| Event | `result` |
|
||||
|-------|----------|
|
||||
| Preview | `previewed` |
|
||||
| Refusal at any stage | `denied` |
|
||||
| Assignment created | `succeeded` |
|
||||
|
||||
`correlation.request_id` carries the request's `correlation_id`, and a
|
||||
successful record's `metadata` carries `assignment_id` and `lease_id`, so an
|
||||
assignment can be traced back to the intent that produced it. The operator's
|
||||
`intent_summary` travels in `metadata` and passes through the standard
|
||||
redaction pass before persistence like every other field.
|
||||
|
||||
## Non-goals
|
||||
|
||||
- No browser-initiated approve or merge, in this phase or any other.
|
||||
- No bypass of allocator exclusive ownership; no self-selection of work.
|
||||
- No auto-start from raw monitoring incidents (#612 stays downstream).
|
||||
@@ -7,6 +7,7 @@ import tempfile
|
||||
import threading
|
||||
import unittest
|
||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from allocator_service import (
|
||||
OUTCOME_ASSIGNED,
|
||||
@@ -15,6 +16,7 @@ from allocator_service import (
|
||||
OUTCOME_PREVIEW,
|
||||
OUTCOME_WAIT,
|
||||
WorkCandidate,
|
||||
_drop_expired_claims,
|
||||
allocate_next_work,
|
||||
candidate_from_dict,
|
||||
classify_skip,
|
||||
@@ -362,5 +364,161 @@ class AllocatorServiceTest(unittest.TestCase):
|
||||
self.assertIn("unavailable", res["reasons"][0].lower())
|
||||
|
||||
|
||||
class SideEffectFreeAllocationTest(unittest.TestCase):
|
||||
"""``side_effect_free`` dry runs write nothing to the control plane (#643).
|
||||
|
||||
A plain ``apply=False`` still called ``upsert_session`` and
|
||||
``expire_stale_leases`` before the apply branch was consulted, so a caller
|
||||
advertising a read-only preview mutated on every call — one unreferenced
|
||||
session row per preview, plus a global lease sweep.
|
||||
"""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self._tmp.cleanup()
|
||||
|
||||
def _alloc(self, **kwargs):
|
||||
defaults = dict(
|
||||
db=self.db,
|
||||
session_id="s-preview",
|
||||
role="author",
|
||||
remote="prgs",
|
||||
org="org",
|
||||
repo="repo",
|
||||
candidates=[
|
||||
WorkCandidate(kind="issue", number=643, labels=("status:ready",))
|
||||
],
|
||||
apply=False,
|
||||
profile_name="prgs-author",
|
||||
username="jcwalker3",
|
||||
)
|
||||
defaults.update(kwargs)
|
||||
return allocate_next_work(**defaults)
|
||||
|
||||
def _session_ids(self) -> set[str]:
|
||||
return {str(r.get("session_id")) for r in self.db.list_sessions()}
|
||||
|
||||
def test_side_effect_free_preview_writes_no_session_row(self):
|
||||
before = self._session_ids()
|
||||
result = self._alloc(side_effect_free=True)
|
||||
self.assertEqual(result["outcome"], OUTCOME_PREVIEW)
|
||||
self.assertEqual(self._session_ids(), before)
|
||||
self.assertNotIn("s-preview", self._session_ids())
|
||||
|
||||
def test_plain_dry_run_still_registers_a_session(self):
|
||||
# The default is unchanged for every existing caller.
|
||||
self._alloc()
|
||||
self.assertIn("s-preview", self._session_ids())
|
||||
|
||||
def test_repeated_previews_do_not_accumulate_rows(self):
|
||||
for index in range(5):
|
||||
self._alloc(side_effect_free=True, session_id=f"s-{index}")
|
||||
self.assertEqual(self._session_ids(), set())
|
||||
|
||||
def test_side_effect_free_does_not_sweep_stale_leases(self):
|
||||
self.db.upsert_session(session_id="owner", role="author", pid=1)
|
||||
assigned = self.db.assign_and_lease(
|
||||
session_id="owner",
|
||||
role="author",
|
||||
remote="prgs",
|
||||
org="org",
|
||||
repo="repo",
|
||||
kind="issue",
|
||||
number=999,
|
||||
lease_ttl_seconds=-60, # already expired
|
||||
)
|
||||
self.assertEqual(assigned.outcome, "assigned")
|
||||
|
||||
self._alloc(side_effect_free=True)
|
||||
|
||||
# The expired row is still 'active' in the DB: nothing swept it.
|
||||
statuses = {
|
||||
r["lease_id"]: r["status"]
|
||||
for r in self.db.list_leases(
|
||||
remote="prgs", org="org", repo="repo",
|
||||
statuses=("active", "expired"),
|
||||
)
|
||||
}
|
||||
self.assertEqual(statuses.get(assigned.lease_id), "active")
|
||||
|
||||
def test_expired_claims_are_filtered_in_memory_so_work_stays_selectable(self):
|
||||
"""The read-only mirror of the sweep: expired claims must not block."""
|
||||
self.db.upsert_session(session_id="owner", role="author", pid=1)
|
||||
self.db.assign_and_lease(
|
||||
session_id="owner",
|
||||
role="author",
|
||||
remote="prgs",
|
||||
org="org",
|
||||
repo="repo",
|
||||
kind="issue",
|
||||
number=643,
|
||||
lease_ttl_seconds=-60, # expired: must not withhold #643
|
||||
)
|
||||
result = self._alloc(side_effect_free=True)
|
||||
self.assertEqual(result["outcome"], OUTCOME_PREVIEW)
|
||||
self.assertEqual(result["selected"]["number"], 643)
|
||||
|
||||
def test_a_live_claim_still_withholds_the_work(self):
|
||||
self.db.upsert_session(session_id="owner", role="author", pid=1)
|
||||
self.db.assign_and_lease(
|
||||
session_id="owner",
|
||||
role="author",
|
||||
remote="prgs",
|
||||
org="org",
|
||||
repo="repo",
|
||||
kind="issue",
|
||||
number=643,
|
||||
lease_ttl_seconds=3600,
|
||||
)
|
||||
result = self._alloc(side_effect_free=True)
|
||||
self.assertNotEqual(result["outcome"], OUTCOME_ASSIGNED)
|
||||
self.assertNotEqual((result.get("selected") or {}).get("number"), 643)
|
||||
|
||||
def test_side_effect_free_with_apply_fails_closed(self):
|
||||
result = self._alloc(side_effect_free=True, apply=True)
|
||||
self.assertFalse(result["success"])
|
||||
self.assertEqual(result["outcome"], OUTCOME_NO_SAFE)
|
||||
self.assertIsNone(result["assignment"])
|
||||
self.assertIn("incompatible with apply", result["reasons"][0])
|
||||
# And it reserved nothing.
|
||||
self.assertEqual(
|
||||
self.db.list_leases(remote="prgs", org="org", repo="repo"), []
|
||||
)
|
||||
|
||||
|
||||
class DropExpiredClaimsTest(unittest.TestCase):
|
||||
"""The in-memory expiry filter behind side-effect-free previews (#643)."""
|
||||
|
||||
def test_unparseable_expiry_is_kept_rather_than_assumed_free(self):
|
||||
claims = {
|
||||
("issue", 1): {"lease_id": "l1", "expires_at": "not-a-date"},
|
||||
("issue", 2): {"lease_id": "l2"},
|
||||
("issue", 3): {"lease_id": "l3", "expires_at": None},
|
||||
}
|
||||
self.assertEqual(_drop_expired_claims(claims), claims)
|
||||
|
||||
def test_expired_dropped_and_future_kept(self):
|
||||
now = datetime(2026, 7, 25, 12, 0, tzinfo=timezone.utc)
|
||||
claims = {
|
||||
("issue", 1): {"expires_at": "2026-07-25T11:59:59+00:00"},
|
||||
("issue", 2): {"expires_at": "2026-07-25T12:00:01+00:00"},
|
||||
("issue", 3): {"expires_at": "2026-07-25T12:00:00+00:00"}, # boundary
|
||||
}
|
||||
kept = _drop_expired_claims(claims, now=now)
|
||||
self.assertEqual(set(kept), {("issue", 2)})
|
||||
|
||||
def test_naive_and_zulu_timestamps_are_treated_as_utc(self):
|
||||
now = datetime(2026, 7, 25, 12, 0, tzinfo=timezone.utc)
|
||||
claims = {
|
||||
("issue", 1): {"expires_at": "2026-07-25T11:00:00"}, # naive, past
|
||||
("issue", 2): {"expires_at": "2026-07-25T13:00:00Z"}, # zulu, future
|
||||
}
|
||||
kept = _drop_expired_claims(claims, now=now)
|
||||
self.assertEqual(set(kept), {("issue", 2)})
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -1,276 +0,0 @@
|
||||
"""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 & 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"]
|
||||
File diff suppressed because it is too large
Load Diff
+111
-20
@@ -72,11 +72,8 @@ from webui.system_health import (
|
||||
snapshot_to_dict as system_health_to_dict,
|
||||
)
|
||||
from webui.system_health_views import render_system_health_page
|
||||
from webui.notifications import (
|
||||
load_notifications_snapshot,
|
||||
snapshot_to_dict as notifications_snapshot_to_dict,
|
||||
)
|
||||
from webui.notification_views import render_notifications_page
|
||||
from webui import request_service
|
||||
from webui.request_views import render_requests_page
|
||||
|
||||
_READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"})
|
||||
_AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"})
|
||||
@@ -744,21 +741,107 @@ 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
|
||||
def _default_request_scope() -> dict[str, str]:
|
||||
"""Resolve remote/org/repo from the project registry for request forms.
|
||||
|
||||
Returns an empty mapping when the registry cannot be read, which makes
|
||||
``parse_request`` reject a request that did not name its own scope rather
|
||||
than letting it default to some other repository.
|
||||
"""
|
||||
from webui.queue_loader import _host_from_url # host normalisation helper
|
||||
|
||||
registry, error = _load_project_registry()
|
||||
if error is not None or not registry.projects:
|
||||
return {}
|
||||
project = registry.projects[0]
|
||||
host = _host_from_url(project.remote_host)
|
||||
return {
|
||||
"remote": _derive_remote(host),
|
||||
"org": project.gitea_owner or "",
|
||||
"repo": project.repo_name or "",
|
||||
}
|
||||
|
||||
|
||||
async def _request_payload(request: Request) -> dict[str, object]:
|
||||
"""Read a request body as JSON or form-encoded. Never raises."""
|
||||
content_type = (request.headers.get("content-type") or "").lower()
|
||||
if "application/json" in content_type:
|
||||
try:
|
||||
body = await request.json()
|
||||
except Exception:
|
||||
return {}
|
||||
return dict(body) if isinstance(body, dict) else {}
|
||||
try:
|
||||
form = await request.form()
|
||||
except Exception:
|
||||
return {}
|
||||
return {key: form[key] for key in form}
|
||||
|
||||
|
||||
async def requests_page(request: Request) -> HTMLResponse:
|
||||
"""Operator request form and intent preview (#643).
|
||||
|
||||
POST here only ever *previews*. Initiation is a separate confirmed call to
|
||||
``/api/v1/requests/apply`` so that submitting this form cannot reserve
|
||||
work as a side effect.
|
||||
"""
|
||||
submitted: dict[str, object] = {}
|
||||
preview = None
|
||||
error = None
|
||||
if request.method == "POST":
|
||||
submitted = await _request_payload(request)
|
||||
work_request, error = request_service.parse_request(
|
||||
submitted, default_scope=_default_request_scope()
|
||||
)
|
||||
if work_request is not None:
|
||||
preview = request_service.preview_request(
|
||||
work_request,
|
||||
principal=resolve_principal(headers=dict(request.headers)),
|
||||
)
|
||||
return HTMLResponse(
|
||||
render_requests_page(
|
||||
preview=preview, error=error, submitted=submitted
|
||||
)
|
||||
)
|
||||
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)
|
||||
async def api_v1_request_preview(request: Request) -> JSONResponse:
|
||||
"""Dry-run authorization and intent preview for a work request (#643)."""
|
||||
payload = await _request_payload(request)
|
||||
work_request, error = request_service.parse_request(
|
||||
payload, default_scope=_default_request_scope()
|
||||
)
|
||||
if work_request is None:
|
||||
return JSONResponse(error.to_dict(), status_code=400)
|
||||
preview = request_service.preview_request(
|
||||
work_request,
|
||||
principal=resolve_principal(headers=dict(request.headers)),
|
||||
)
|
||||
return JSONResponse(
|
||||
preview.to_dict(), status_code=200 if preview.authorized else 403
|
||||
)
|
||||
|
||||
|
||||
async def api_v1_request_apply(request: Request) -> JSONResponse:
|
||||
"""Initiate a previewed work request through the allocator (#643).
|
||||
|
||||
Fail-closed at every step: unauthorized, unconfirmed, not-next-safe, and
|
||||
already-claimed all return without attempting an assignment.
|
||||
"""
|
||||
payload = await _request_payload(request)
|
||||
work_request, error = request_service.parse_request(
|
||||
payload, default_scope=_default_request_scope()
|
||||
)
|
||||
if work_request is None:
|
||||
return JSONResponse(error.to_dict(), status_code=400)
|
||||
confirm = _truthy_flag(str(payload.get("confirm") or ""))
|
||||
result = request_service.apply_request(
|
||||
work_request,
|
||||
principal=resolve_principal(headers=dict(request.headers)),
|
||||
confirm=confirm,
|
||||
)
|
||||
status = int(result.pop("status_code", 403))
|
||||
return JSONResponse(result, status_code=status)
|
||||
|
||||
|
||||
async def method_not_allowed(request: Request, _exc: Exception) -> Response:
|
||||
@@ -789,9 +872,6 @@ 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"]),
|
||||
@@ -831,6 +911,17 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
|
||||
api_action_attempt,
|
||||
methods=["POST"],
|
||||
),
|
||||
Route("/requests", requests_page, methods=["GET", "POST"]),
|
||||
Route(
|
||||
"/api/v1/requests/preview",
|
||||
api_v1_request_preview,
|
||||
methods=["POST"],
|
||||
),
|
||||
Route(
|
||||
"/api/v1/requests/apply",
|
||||
api_v1_request_apply,
|
||||
methods=["POST"],
|
||||
),
|
||||
Route("/api/leases", api_leases, methods=["GET"]),
|
||||
Route("/api/v1/inventory", api_inventory, methods=["GET"]),
|
||||
Route(
|
||||
|
||||
+71
-8
@@ -115,6 +115,12 @@ class ConsoleAction:
|
||||
break_glass: bool
|
||||
phase: int
|
||||
summary: str
|
||||
# Opt-in switch for an action whose execution path is genuinely wired
|
||||
# ahead of its phase becoming globally active (#643). Naming a variable
|
||||
# here enables nothing on its own: the variable must also be set in the
|
||||
# environment. An action that leaves this ``None`` can only execute once
|
||||
# ACTIVE_PHASE reaches its phase, exactly as before.
|
||||
execution_env_flag: str | None = None
|
||||
|
||||
@property
|
||||
def mcp_permission(self) -> str:
|
||||
@@ -277,6 +283,27 @@ _ACTION_SPECS: tuple[ConsoleAction, ...] = (
|
||||
phase=2,
|
||||
summary="Restart one MCP namespace via the host supervisor.",
|
||||
),
|
||||
# #643: submit a work request — desired role, issue/PR, intent — and let
|
||||
# the allocator reserve it. This is the one Phase 2 action whose execution
|
||||
# path is actually implemented (``webui.request_service``), so it carries
|
||||
# the opt-in flag; it stays denied until an operator sets that variable.
|
||||
# Authority is operator-class because the outcome is a claim, not a Gitea
|
||||
# verdict: initiating reviewer or merger *work* does not grant the right
|
||||
# to approve or merge, which stays with the MCP role profile.
|
||||
ConsoleAction(
|
||||
action_id="initiate_workflow",
|
||||
task_key="allocate_next_work",
|
||||
action_class=CLASS_WRITE,
|
||||
minimum_role=OPERATOR,
|
||||
requires_confirmation=True,
|
||||
dual_control=False,
|
||||
break_glass=False,
|
||||
phase=2,
|
||||
summary=(
|
||||
"Preview and initiate allocator-owned workflow work for a role."
|
||||
),
|
||||
execution_env_flag="WEBUI_REQUESTS_EXECUTION",
|
||||
),
|
||||
)
|
||||
|
||||
ACTIONS: dict[str, ConsoleAction] = {a.action_id: a for a in _ACTION_SPECS}
|
||||
@@ -430,6 +457,33 @@ ALLOW_PREVIEW = "allowed_preview_only"
|
||||
# gated on this model landing; nothing here enables it.
|
||||
ACTIVE_PHASE = 1
|
||||
|
||||
_TRUTHY = frozenset({"1", "true", "yes", "on"})
|
||||
|
||||
|
||||
def execution_wired(
|
||||
action: ConsoleAction | None, env: dict[str, str] | None = None
|
||||
) -> bool:
|
||||
"""Whether *action* has a live execution path right now.
|
||||
|
||||
Two ways to be wired, and only two. The action's phase is active, or the
|
||||
action declares an opt-in environment variable *and* that variable is set.
|
||||
Everything else — including every action that never declares a flag — is
|
||||
unwired, so the default across the registry stays deny.
|
||||
|
||||
Bumping ``ACTIVE_PHASE`` would enable execution for every action of that
|
||||
phase at once. The per-action flag exists so a single implemented action
|
||||
can go live without dragging its unimplemented phase-mates with it.
|
||||
"""
|
||||
if action is None:
|
||||
return False
|
||||
if action.phase <= ACTIVE_PHASE:
|
||||
return True
|
||||
flag = (action.execution_env_flag or "").strip()
|
||||
if not flag:
|
||||
return False
|
||||
source = env if env is not None else os.environ
|
||||
return (source.get(flag) or "").strip().lower() in _TRUTHY
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class AuthorizationDecision:
|
||||
@@ -469,16 +523,19 @@ def authorize(
|
||||
principal: Principal | None = None,
|
||||
*,
|
||||
for_execution: bool = False,
|
||||
env: dict[str, str] | None = None,
|
||||
) -> AuthorizationDecision:
|
||||
"""Decide whether *principal* may invoke *action_id*. Deny by default.
|
||||
|
||||
``for_execution`` distinguishes a read-only preview from a real invocation.
|
||||
Even an allowed decision reports ``execution_enabled=False`` while the
|
||||
console is in Phase 1, so no caller can read an allow as permission to
|
||||
mutate.
|
||||
``execution_enabled`` reports whether the action has a live execution path
|
||||
at all (:func:`execution_wired`) — for every action without an explicit
|
||||
opt-in flag that stays ``False`` while the console is in Phase 1, so no
|
||||
caller can read an allow as permission to mutate.
|
||||
"""
|
||||
who = principal if principal is not None else ANONYMOUS
|
||||
action = get_action(action_id)
|
||||
wired = execution_wired(action, env)
|
||||
|
||||
if action is None:
|
||||
return AuthorizationDecision(
|
||||
@@ -497,7 +554,7 @@ def authorize(
|
||||
"requires_confirmation": action.requires_confirmation,
|
||||
"dual_control": action.dual_control,
|
||||
"break_glass": action.break_glass,
|
||||
"execution_enabled": False,
|
||||
"execution_enabled": wired,
|
||||
}
|
||||
|
||||
if not who.authenticated:
|
||||
@@ -530,13 +587,19 @@ def authorize(
|
||||
**base,
|
||||
)
|
||||
|
||||
if for_execution and action.phase > ACTIVE_PHASE:
|
||||
if for_execution and not wired:
|
||||
return AuthorizationDecision(
|
||||
allowed=False,
|
||||
reason_code=DENY_PHASE_NOT_ACTIVE,
|
||||
detail=(
|
||||
f"Action {action_id!r} belongs to phase {action.phase}; the "
|
||||
f"console is in phase {ACTIVE_PHASE}. Execution is not wired."
|
||||
f"console is in phase {ACTIVE_PHASE}"
|
||||
+ (
|
||||
f" and {action.execution_env_flag} is not set"
|
||||
if action.execution_env_flag
|
||||
else ""
|
||||
)
|
||||
+ ". Execution is not wired."
|
||||
),
|
||||
**base,
|
||||
)
|
||||
@@ -545,8 +608,8 @@ def authorize(
|
||||
allowed=True,
|
||||
reason_code=ALLOW_PREVIEW,
|
||||
detail=(
|
||||
"Principal holds the required role. Preview only — execution "
|
||||
"remains disabled until the Phase 2 action framework ships."
|
||||
"Principal holds the required role. Execution proceeds only for an "
|
||||
"action with a wired execution path; everything else is preview."
|
||||
),
|
||||
**base,
|
||||
)
|
||||
|
||||
+1
-1
@@ -45,7 +45,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
|
||||
NavItem("/queue", "Queue"),
|
||||
NavItem("/leases", "Leases"),
|
||||
NavItem("/actions", "Actions"),
|
||||
NavItem("/notifications", "Notifications"),
|
||||
NavItem("/requests", "Requests"),
|
||||
)),
|
||||
NavGroup("Runtime/Sessions", (
|
||||
NavItem("/runtime", "Runtime health"),
|
||||
|
||||
@@ -1,158 +0,0 @@
|
||||
"""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)
|
||||
@@ -1,480 +0,0 @@
|
||||
"""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
|
||||
|
||||
import os
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Callable, Sequence
|
||||
|
||||
from webui import console_redaction
|
||||
from webui.project_registry import find_project, 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).
|
||||
"""
|
||||
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}
|
||||
or "hard stop" in summary.lower()
|
||||
or "unauthorized" in summary.lower()
|
||||
or "irrecoverable" in summary.lower()
|
||||
):
|
||||
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 probe_err in 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}",
|
||||
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 collision in lease_snap.duplicate_prs:
|
||||
att_cls, req_human = classify_attention_event(
|
||||
CATEGORY_BLOCKER,
|
||||
f"Duplicate PR Collision ({collision.kind})",
|
||||
collision.message,
|
||||
is_blocker=True,
|
||||
)
|
||||
items.append(
|
||||
NotificationItem(
|
||||
id=f"notif-collision-{collision.issue_number or 0}",
|
||||
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_err = queue_snap.fetch_error or lease_snap.fetch_error or getattr(health_snap, "probe_errors", None)
|
||||
if isinstance(fetch_err, (tuple, list)):
|
||||
fetch_err = "; ".join(fetch_err) if fetch_err else None
|
||||
|
||||
return NotificationSnapshot(
|
||||
project_id=project.id,
|
||||
repo_label=f"{project.gitea_owner}/{project.repo_name}",
|
||||
items=tuple(items),
|
||||
human_required_count=human_req_count,
|
||||
operator_count=operator_count,
|
||||
routine_count=routine_count,
|
||||
total_count=len(items),
|
||||
fetch_error=fetch_err,
|
||||
)
|
||||
|
||||
|
||||
def snapshot_to_dict(snapshot: NotificationSnapshot) -> dict[str, Any]:
|
||||
"""JSON-serializable export for /api/v1/notifications."""
|
||||
return snapshot.as_dict()
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,164 @@
|
||||
"""HTML views for the operator request surface (#643).
|
||||
|
||||
The form is deliberately a *preview* form. It has no initiate button, because
|
||||
initiating requires a confirmed POST to ``/api/v1/requests/apply`` and a stray
|
||||
form submission must not be able to produce one by accident.
|
||||
|
||||
Nothing rendered here is trusted input: every interpolated value is escaped,
|
||||
and the page renders only values the service already produced rather than
|
||||
echoing a raw request body back.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import html
|
||||
import json
|
||||
from typing import Any
|
||||
|
||||
from webui.layout import render_page
|
||||
from webui.request_service import (
|
||||
REQUESTABLE_ROLES,
|
||||
WORK_KINDS,
|
||||
RequestError,
|
||||
RequestPreview,
|
||||
)
|
||||
|
||||
REQUESTS_PATH = "/requests"
|
||||
PREVIEW_API_PATH = "/api/v1/requests/preview"
|
||||
APPLY_API_PATH = "/api/v1/requests/apply"
|
||||
|
||||
|
||||
def _escape(text: Any) -> str:
|
||||
return html.escape(str(text if text is not None else ""), quote=True)
|
||||
|
||||
|
||||
REQUEST_PAGE_STYLES = """
|
||||
<style>
|
||||
.request-form { display: grid; gap: 0.75rem; max-width: 44rem; }
|
||||
.request-form label { display: grid; gap: 0.25rem; font-size: 0.9rem; }
|
||||
.request-check { margin: 0.35rem 0; }
|
||||
.request-check .verdict-ok { color: var(--accent); }
|
||||
.request-check .verdict-fail { color: #d14; }
|
||||
.request-prohibited code { margin-right: 0.4rem; }
|
||||
</style>
|
||||
"""
|
||||
|
||||
|
||||
def _options(values: tuple[str, ...], selected: Any) -> str:
|
||||
return "".join(
|
||||
f"<option value='{_escape(value)}'"
|
||||
+ (" selected" if selected == value else "")
|
||||
+ f">{_escape(value)}</option>"
|
||||
for value in values
|
||||
)
|
||||
|
||||
|
||||
def _form(values: dict[str, Any] | None = None) -> str:
|
||||
current = dict(values or {})
|
||||
number = current.get("work_number")
|
||||
return (
|
||||
f"<form class='request-form' method='post' action='{REQUESTS_PATH}'>"
|
||||
"<label>Desired role<select name='desired_role'>"
|
||||
f"{_options(REQUESTABLE_ROLES, current.get('desired_role'))}"
|
||||
"</select></label>"
|
||||
"<label>Work kind<select name='work_kind'>"
|
||||
f"{_options(WORK_KINDS, current.get('work_kind'))}"
|
||||
"</select></label>"
|
||||
"<label>Issue or PR number"
|
||||
"<input type='number' name='work_number' min='1' required "
|
||||
f"value='{_escape(number) if number else ''}'></label>"
|
||||
"<label>Intent summary"
|
||||
"<input type='text' name='intent_summary' maxlength='500' required "
|
||||
f"value='{_escape(current.get('intent_summary'))}'></label>"
|
||||
"<label>Expected head SHA <span class='muted'>(PR work only)</span>"
|
||||
"<input type='text' name='expected_head_sha' "
|
||||
f"value='{_escape(current.get('expected_head_sha'))}'></label>"
|
||||
"<button type='submit' class='copy-btn'>Preview request</button>"
|
||||
"<p class='muted meta'>Preview is read-only and creates no assignment. "
|
||||
f"Initiating requires a confirmed POST to <code>{APPLY_API_PATH}</code>."
|
||||
"</p>"
|
||||
"</form>"
|
||||
)
|
||||
|
||||
|
||||
def _checks_block(preview: RequestPreview) -> str:
|
||||
rows = []
|
||||
for check in preview.checks:
|
||||
verdict = "PASS" if check.ok else "FAIL"
|
||||
css = "verdict-ok" if check.ok else "verdict-fail"
|
||||
rows.append(
|
||||
"<li class='request-check'>"
|
||||
f"<span class='{css}'><strong>{verdict}</strong></span> "
|
||||
f"<code>{_escape(check.name)}</code> — {_escape(check.detail)} "
|
||||
f"<span class='muted meta'>({_escape(check.reason_code)})</span>"
|
||||
"</li>"
|
||||
)
|
||||
return "<ul>" + "".join(rows) + "</ul>"
|
||||
|
||||
|
||||
def _preview_block(preview: RequestPreview) -> str:
|
||||
verdict = "AUTHORIZED" if preview.authorized else "DENIED"
|
||||
prohibited = "".join(
|
||||
f"<code>{_escape(action)}</code>" for action in preview.prohibited_actions
|
||||
)
|
||||
request = preview.request
|
||||
evidence = json.dumps(preview.allocator_evidence, indent=2, default=str)
|
||||
return (
|
||||
"<h3>Intent preview</h3>"
|
||||
f"<p><strong>{verdict}</strong> — {_escape(preview.detail)}</p>"
|
||||
"<p class='meta'>"
|
||||
f"Role <code>{_escape(request.desired_role)}</code> · "
|
||||
f"{_escape(request.work_kind)} <code>{_escape(request.display_ref)}</code>"
|
||||
f" · profile <code>{_escape(preview.required_profile)}</code> · "
|
||||
f"namespace <code>{_escape(preview.required_namespace)}</code> · "
|
||||
f"permission <code>{_escape(preview.required_permission)}</code>"
|
||||
"</p>"
|
||||
f"<p>Intent: {_escape(request.intent_summary)}</p>"
|
||||
f"{_checks_block(preview)}"
|
||||
f"<p><strong>Next safe action:</strong> "
|
||||
f"{_escape(preview.next_safe_action)}</p>"
|
||||
"<p class='request-prohibited'><strong>Prohibited for this role:</strong> "
|
||||
+ (prohibited or "<span class='muted'>none declared</span>")
|
||||
+ "</p>"
|
||||
"<p class='muted meta'>Correlation id "
|
||||
f"<code>{_escape(preview.correlation_id)}</code></p>"
|
||||
"<details><summary>Allocator evidence</summary>"
|
||||
f"<pre class='prompt-text'>{_escape(evidence)}</pre>"
|
||||
"</details>"
|
||||
)
|
||||
|
||||
|
||||
def _error_block(error: RequestError) -> str:
|
||||
field = (
|
||||
f"<p class='meta'>Field: <code>{_escape(error.field_name)}</code></p>"
|
||||
if error.field_name
|
||||
else ""
|
||||
)
|
||||
return (
|
||||
"<h3>Request rejected</h3>"
|
||||
f"<p><strong>{_escape(error.reason_code)}</strong> — "
|
||||
f"{_escape(error.detail)}</p>{field}"
|
||||
)
|
||||
|
||||
|
||||
def render_requests_page(
|
||||
*,
|
||||
preview: RequestPreview | None = None,
|
||||
error: RequestError | None = None,
|
||||
submitted: dict[str, Any] | None = None,
|
||||
) -> str:
|
||||
"""Render the request form, plus a preview or rejection when one exists."""
|
||||
body = (
|
||||
"<h2>Requests</h2>"
|
||||
"<p>Submit a work request — desired role, issue or PR, and intent — "
|
||||
"and see whether it would be authorized before anything is reserved. "
|
||||
"Initiation goes through the allocator (#600/#613); this console never "
|
||||
"self-selects work, never approves, and never merges.</p>"
|
||||
+ _form(submitted)
|
||||
+ (_error_block(error) if error is not None else "")
|
||||
+ (_preview_block(preview) if preview is not None else "")
|
||||
+ f"<p class='meta'><a href='{PREVIEW_API_PATH}'>Preview API</a> · "
|
||||
"<a href='/api/console/security-model'>RBAC model</a></p>"
|
||||
+ REQUEST_PAGE_STYLES
|
||||
)
|
||||
return render_page(title="Requests", body_html=body)
|
||||
@@ -201,6 +201,16 @@ def _candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate
|
||||
return candidates
|
||||
|
||||
|
||||
def candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate]:
|
||||
"""Public alias for :func:`_candidates_from_queue_snapshot` (#643).
|
||||
|
||||
The request-initiation service ranks the same candidate set this view
|
||||
renders, so both must agree on how a queue row becomes a candidate. One
|
||||
construction, two callers — not two that can drift apart.
|
||||
"""
|
||||
return _candidates_from_queue_snapshot(q_snap)
|
||||
|
||||
|
||||
def _claim_lease_records(inventory: dict[str, Any] | None) -> list[dict[str, Any]]:
|
||||
"""Normalize ``build_claim_inventory`` entries into lease records.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user