Address PR #885 REQUEST_CHANGES: full head_sha pins from queue signals, reviewer leases keyed by pr_number only, claim inventory via entries, live-path fixture tests, and traffic state vocabulary docs. Closes #640 (re-review at new head) Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
450 lines
16 KiB
Python
450 lines
16 KiB
Python
"""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()
|