"""Traffic-control view loader for Phase 1 operator web console (#640). Combines queue snapshots, inventory leases, dependency graph classifications, and workflow dashboard rules to deliver full traffic-control visibility: runnable, leased (in-progress), blocked (dependency/lock), needs-controller, and terminal-complete candidates. """ from __future__ import annotations import os from dataclasses import dataclass from typing import Any, Callable, Sequence from webui.project_registry import find_project, load_registry from webui.queue_loader import load_queue_snapshot, QueueSnapshot from webui.lease_loader import load_lease_snapshot, LeaseSnapshot from workflow_dashboard import ( DashboardSnapshot, QueueEntry, RoleNextAction, build_workflow_dashboard, DASHBOARD_ROLES, ) from allocator_service import WorkCandidate @dataclass(frozen=True) class TrafficItem: kind: str # "issue" or "pr" number: int title: str traffic_state: str # "runnable", "leased", "blocked", "needs_controller", "terminal_complete" expected_role: str safe_for_roles: tuple[str, ...] badges: tuple[str, ...] block_reason: str | None = None lease_info: dict[str, Any] | None = None head_sha: str | None = None @property def is_safe(self) -> bool: return self.block_reason is None and bool(self.safe_for_roles) def as_dict(self) -> dict[str, Any]: return { "kind": self.kind, "number": self.number, "title": self.title, "traffic_state": self.traffic_state, "expected_role": self.expected_role, "safe_for_roles": list(self.safe_for_roles), "badges": list(self.badges), "block_reason": self.block_reason, "lease_info": self.lease_info, "head_sha": self.head_sha, "is_safe": self.is_safe, } @dataclass(frozen=True) class TrafficSnapshot: project_id: str repo_label: str runnable: tuple[TrafficItem, ...] leased: tuple[TrafficItem, ...] blocked: tuple[TrafficItem, ...] needs_controller: tuple[TrafficItem, ...] terminal_complete: tuple[TrafficItem, ...] next_roles: tuple[dict[str, Any], ...] fetch_error: str | None = None inventory_complete: bool = True def as_dict(self) -> dict[str, Any]: return { "project_id": self.project_id, "repo_label": self.repo_label, "runnable": [i.as_dict() for i in self.runnable], "leased": [i.as_dict() for i in self.leased], "blocked": [i.as_dict() for i in self.blocked], "needs_controller": [i.as_dict() for i in self.needs_controller], "terminal_complete": [i.as_dict() for i in self.terminal_complete], "next_roles": list(self.next_roles), "fetch_error": self.fetch_error, "inventory_complete": self.inventory_complete, } def _classify_traffic_item( entry: QueueEntry, *, lease_info: dict[str, Any] | None = None, ) -> TrafficItem: """Classify a QueueEntry into a TrafficItem with explicit traffic state.""" badges = list(entry.badges) block_reason = entry.block_reason expected_role = entry.expected_role entry_is_safe = entry.block_reason is None and bool(entry.safe_for_roles) if lease_info is not None or "in-progress" in badges or "claimed" in badges: state = "leased" elif expected_role == "reconciler" or "terminal-lock" in badges: state = "terminal_complete" elif expected_role == "controller" or "contaminated" in badges or "needs-controller" in badges: state = "needs_controller" elif block_reason is not None or "blocked" in badges or "dependency-unmet" in badges or "blocked-by-terminal" in badges: state = "blocked" elif entry_is_safe: state = "runnable" else: state = "needs_controller" return TrafficItem( kind=entry.kind, number=entry.number, title=entry.title, traffic_state=state, expected_role=expected_role, safe_for_roles=entry.safe_for_roles, badges=tuple(badges), block_reason=block_reason, lease_info=lease_info, head_sha=entry.head_sha, ) def load_traffic_snapshot( *, candidates: Sequence[WorkCandidate] | None = None, leases: Sequence[dict[str, Any]] | None = None, terminal_pr: int | None = None, fetch_queue_snapshot: Callable[[], QueueSnapshot] | None = None, fetch_lease_snapshot: Callable[[], LeaseSnapshot] | None = None, project_id: str = "gitea-tools", ) -> TrafficSnapshot: """Load and compute the traffic-control snapshot.""" try: reg = load_registry() proj = find_project(reg, project_id) repo_label = proj.remote_repo if proj else "Scaled-Tech-Consulting/Gitea-Tools" except Exception: repo_label = "Scaled-Tech-Consulting/Gitea-Tools" # Injected candidates path (pure unit testing) if candidates is not None: dashboard = build_workflow_dashboard( candidates=candidates, leases=leases, terminal_pr=terminal_pr, inventory_complete=True, ) return _build_traffic_snapshot_from_dashboard( project_id=project_id, repo_label=repo_label, dashboard=dashboard, leases=leases or (), ) # Live snapshot loading q_loader = fetch_queue_snapshot or load_queue_snapshot l_loader = fetch_lease_snapshot or load_lease_snapshot try: q_snap = q_loader() l_snap = l_loader() except Exception as exc: # noqa: BLE001 return TrafficSnapshot( project_id=project_id, repo_label=repo_label, runnable=(), leased=(), blocked=(), needs_controller=(), terminal_complete=(), next_roles=(), fetch_error=f"Failed to load traffic state: {exc}", inventory_complete=False, ) if q_snap.fetch_error or l_snap.fetch_error: err = q_snap.fetch_error or l_snap.fetch_error return TrafficSnapshot( project_id=project_id, repo_label=repo_label, runnable=(), leased=(), blocked=(), needs_controller=(), terminal_complete=(), next_roles=(), fetch_error=err, inventory_complete=False, ) # Build WorkCandidates from live queue snapshot candidate_list: list[WorkCandidate] = [] for pr in q_snap.prs: linked = int(pr.extra["linked_issue"]) if pr.extra.get("linked_issue") and pr.extra["linked_issue"].isdigit() else None candidate_list.append( WorkCandidate( kind="pr", number=pr.number, state="open", labels=(), title=pr.title, priority=10 if "request-changes" in pr.badges else 5, request_changes_current_head="request-changes" in pr.badges, approval_on_current_head="approved" in pr.badges or "merge-ready" in pr.badges, mergeable="blocked" not in pr.badges, blocked="blocked" in pr.badges, linked_issue_number=linked, ) ) for issue in q_snap.issues: labels = [b for b in issue.badges if b.startswith("status:") or b in ("discussion", "blocked", "ready", "in-progress")] candidate_list.append( WorkCandidate( kind="issue", number=issue.number, state="open", labels=tuple(labels), title=issue.title, priority=20 if "status:ready" in labels else 10, blocked="blocked" in labels or "status:blocked" in labels, ) ) raw_leases: list[dict[str, Any]] = [] if l_snap.claim_inventory and "active_claims" in l_snap.claim_inventory: raw_leases.extend(l_snap.claim_inventory["active_claims"]) for r_lease in l_snap.reviewer_leases: raw_leases.append(r_lease) dashboard = build_workflow_dashboard( candidates=candidate_list, leases=raw_leases, inventory_complete=q_snap.pr_pagination.inventory_complete if q_snap.pr_pagination else True, ) return _build_traffic_snapshot_from_dashboard( project_id=project_id, repo_label=repo_label, dashboard=dashboard, leases=raw_leases, ) def _build_traffic_snapshot_from_dashboard( *, project_id: str, repo_label: str, dashboard: DashboardSnapshot, leases: Sequence[dict[str, Any]], ) -> TrafficSnapshot: """Classify dashboard entries into the 5 traffic state buckets.""" all_entries = dashboard.open_prs + dashboard.open_issues # Map leased work numbers lease_map: dict[tuple[str, int], dict[str, Any]] = {} for lease in leases: if isinstance(lease, dict): kind = str(lease.get("kind") or lease.get("work_kind") or "issue").lower() num = lease.get("number") or lease.get("work_number") or lease.get("issue_number") or lease.get("pr_number") if num is not None: try: lease_map[(kind, int(num))] = lease except (ValueError, TypeError): pass runnable: list[TrafficItem] = [] leased: list[TrafficItem] = [] blocked: list[TrafficItem] = [] needs_controller: list[TrafficItem] = [] terminal_complete: list[TrafficItem] = [] for entry in all_entries: l_info = lease_map.get((entry.kind, entry.number)) item = _classify_traffic_item(entry, lease_info=l_info) if item.traffic_state == "leased": leased.append(item) elif item.traffic_state == "terminal_complete": terminal_complete.append(item) elif item.traffic_state == "blocked": blocked.append(item) elif item.traffic_state == "needs_controller": needs_controller.append(item) else: runnable.append(item) next_roles = [dashboard.next_safe_by_role[r].as_dict() for r in DASHBOARD_ROLES if r in dashboard.next_safe_by_role] return TrafficSnapshot( project_id=project_id, repo_label=repo_label, runnable=tuple(runnable), leased=tuple(leased), blocked=tuple(blocked), needs_controller=tuple(needs_controller), terminal_complete=tuple(terminal_complete), next_roles=tuple(next_roles), fetch_error=None, inventory_complete=dashboard.inventory_complete, ) def snapshot_to_dict(snapshot: TrafficSnapshot) -> dict[str, Any]: return snapshot.as_dict()