"""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 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) # Lease state is checked first: an item that is both leased and blocked is # reported as leased. That is safe by construction — a leased item is never # placed in the runnable lane — and it keeps the operator's attention on the # session that currently owns the work. The blocker text still renders. 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 or "status:blocked" 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, ) # Claim statuses from ``issue_claim_heartbeat.build_claim_inventory`` that mean # a live worker currently holds the issue. Everything else (``stale``, # ``phantom``, ``reclaimable``, ``not_claimed``) is reported through the # dashboard's stale-lease channel and is never rendered as an active lease. _ACTIVE_CLAIM_STATUSES = frozenset({"active", "awaiting_review"}) # Statuses that positively mean "not an active lease" for any lease record. _INACTIVE_LEASE_STATUSES = frozenset( {"expired", "stale", "released", "moot", "reclaimable", "phantom", "not_claimed"} ) def _candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate]: """Build allocator candidates from the queue loader's authoritative signals. Display badges (``blocked``/``claimed``/``duplicate``/``stale``/ ``in-review``/``open``) are rendering hints, not routing state, so nothing here branches on them. Every routing field comes from ``QueueItem.signals`` — the raw Gitea payload values. The queue loader reads ``/pulls`` and ``/issues`` only; it never fetches review verdicts. ``request_changes_current_head`` / ``approval_on_current_head`` are therefore left at their fail-safe ``False`` rather than being guessed from badges: an unproven approval must never route a PR to the merger. """ candidates: list[WorkCandidate] = [] for pr in q_snap.prs: signals = pr.signals or {} head_sha = str(signals.get("head_sha") or "").strip() mergeable = signals.get("mergeable") labels = tuple(str(x) for x in (signals.get("labels") or ())) candidates.append( WorkCandidate( kind="pr", number=pr.number, state="open", labels=labels, title=pr.title, # Full 40-char SHA from head.sha — never the 12-char display value. head_sha=head_sha or None, priority=5, mergeable=mergeable is True, blocked=mergeable is False or "status:blocked" in labels, ) ) for issue in q_snap.issues: signals = issue.signals or {} labels = tuple(str(x) for x in (signals.get("labels") or ())) lowered = {label.lower() for label in labels} candidates.append( WorkCandidate( kind="issue", number=issue.number, state="open", labels=labels, title=issue.title, priority=20 if "status:ready" in lowered else 10, blocked="status:blocked" in lowered, # A live claim by another session is not this session's work. already_claimed_elsewhere="status:in-progress" in lowered, ) ) return candidates def _claim_lease_records(inventory: dict[str, Any] | None) -> list[dict[str, Any]]: """Normalize ``build_claim_inventory`` entries into lease records. The inventory contract is ``{"entries", "counts", "heartbeat_lease_minutes", "reclaim_after_minutes", "in_progress_total"}``. Each entry is keyed by ``issue_number``; the subject kind is therefore always ``issue``. """ entries = (inventory or {}).get("entries") or () records: list[dict[str, Any]] = [] for entry in entries: if not isinstance(entry, dict): continue number = entry.get("issue_number") if number is None: continue try: number_int = int(number) except (TypeError, ValueError): continue heartbeat = entry.get("latest_heartbeat") or {} record = dict(entry) record.update( { "kind": "issue", "number": number_int, "role": "author", "lease_source": "issue-claim-heartbeat", } ) if isinstance(heartbeat, dict): if heartbeat.get("session_id") and not record.get("session_id"): record["session_id"] = heartbeat.get("session_id") if heartbeat.get("author") and not record.get("author"): record["author"] = heartbeat.get("author") records.append(record) return records def _lease_subject(lease: dict[str, Any]) -> tuple[str, int] | None: """Return the ``(kind, number)`` a lease record actually covers. Fails closed: a record that does not identify exactly one subject is dropped rather than attributed to a guessed work item (#640 — never invent a lease, and never attach a PR lease to a same-numbered issue). """ kind = str(lease.get("kind") or lease.get("work_kind") or "").strip().lower() pr_number = lease.get("pr_number") issue_number = lease.get("issue_number") if kind not in ("pr", "issue"): if pr_number is not None and issue_number is None: kind = "pr" elif issue_number is not None and pr_number is None: kind = "issue" else: return None number = lease.get("number") if number is None: number = lease.get("work_number") if number is None: number = pr_number if kind == "pr" else issue_number if number is None: return None try: return kind, int(number) except (TypeError, ValueError): return None def _is_active_lease(lease: dict[str, Any]) -> bool: """True when the record proves a worker currently holds the item.""" if lease.get("stale") or lease.get("expired"): return False status = str(lease.get("status") or lease.get("lease_status") or "").strip().lower() if status in _INACTIVE_LEASE_STATUSES: return False if lease.get("lease_source") == "issue-claim-heartbeat": return status in _ACTIVE_CLAIM_STATUSES return True 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, ) candidate_list = _candidates_from_queue_snapshot(q_snap) raw_leases: list[dict[str, Any]] = _claim_lease_records(l_snap.claim_inventory) for r_lease in l_snap.reviewer_leases or (): if not isinstance(r_lease, dict): continue # Always pin reviewer leases to the PR subject, even if a linked # issue_number is present on the marker (#640 B2). normalized = dict(r_lease) subject = normalized.get("pr_number") or normalized.get("number") if subject is None: continue try: pr_num = int(subject) except (TypeError, ValueError): continue normalized["kind"] = "pr" normalized["number"] = pr_num normalized["pr_number"] = pr_num normalized.setdefault("role", "reviewer") raw_leases.append(normalized) 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 each active lease onto the exact work item it covers. Records whose # subject cannot be determined, and claims that are stale/phantom/ # reclaimable, are deliberately dropped instead of guessed. lease_map: dict[tuple[str, int], dict[str, Any]] = {} for lease in leases: if not isinstance(lease, dict) or not _is_active_lease(lease): continue subject = _lease_subject(lease) if subject is not None: lease_map[subject] = lease 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()