"""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()