Merge commit '578c44b685a7ff5b01006c5e398bfac9863e0d8d' into feat/issue-661-drain-proof-hard-gate

This commit is contained in:
2026-07-24 22:37:49 -04:00
8 changed files with 1115 additions and 4 deletions
+31
View File
@@ -57,6 +57,8 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| `/system-health` | System-health dashboard — readiness, version/uptime, dependencies, MCP namespaces, stale-runtime parity (#639) | | `/system-health` | System-health dashboard — readiness, version/uptime, dependencies, MCP namespaces, stale-runtime parity (#639) |
| `/queue` | Live PR and issue queue dashboard (#429) | | `/queue` | Live PR and issue queue dashboard (#429) |
| `/api/queue` | JSON queue export with pagination metadata | | `/api/queue` | JSON queue export with pagination metadata |
| `/traffic` | Workflow traffic-control view — runnable, leased, blocked, needs-controller, terminal-complete (#640) |
| `/api/traffic` | JSON traffic-control export with state classifications and next safe role actions |
| `/projects` | Project registry list with status and onboarding progress (#427, #635) | | `/projects` | Project registry list with status and onboarding progress (#427, #635) |
| `/projects/{id}` | Project detail + onboarding checklist | | `/projects/{id}` | Project detail + onboarding checklist |
| `/api/v1/projects` | Versioned JSON registry export (#635) | | `/api/v1/projects` | Versioned JSON registry export (#635) |
@@ -85,6 +87,35 @@ Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for `read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
local validator preview only (no Gitea mutations, no server-side storage). local validator preview only (no Gitea mutations, no server-side storage).
### Traffic-control state vocabulary (#640)
The traffic view classifies each open issue/PR into exactly one bucket:
| Bucket | Meaning | Operator implication |
|--------|---------|----------------------|
| **runnable** | No active lease, no block reason, safe for its expected role | Next role may start work |
| **leased** | Active author claim or reviewer PR lease | Do not stomp; wait or adopt via role tools |
| **blocked** | Dependency, missing head pin, conflict, or unmet dependency | Author remediation first |
| **needs_controller** | Contaminated, controller-only diagnosis, or `status:blocked` | Controller only |
| **terminal_complete** | Reconciler / terminal-lock territory | Reconciler cleanup path |
`status:blocked` items route to **needs_controller**, not **blocked**:
`expected_role_for_candidate` sends them to the controller, and the blocker
reason renders in either bucket.
**Live path contracts (do not invent):**
- PR head pins come from `QueueItem.signals["head_sha"]` (full SHA). Display
`extra["head_sha"]` is truncated and must never be used for routing.
- Reviewer leases are keyed as `(pr, pr_number)` only — never via a linked
`issue_number` on the same lease marker.
- Issue claims come from `claim_inventory["entries"]`
(`issue_claim_heartbeat.build_claim_inventory`). There is no `active_claims`
key.
- Queue display badges are only: `blocked`, `claimed`, `duplicate`, `stale`,
`in-review`, `open`. Review verdicts (`request-changes`, `approved`) are
**not** queue badges; traffic does not invent them from the queue loader.
## System health API (#634) ## System health API (#634)
`GET /api/v1/system/health` is the structured, read-only health surface for `GET /api/v1/system/health` is the structured, read-only health surface for
+408
View File
@@ -0,0 +1,408 @@
"""Tests for web UI workflow traffic-control view (#640)."""
import sys
import unittest
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient
from webui.app import create_app
from webui.traffic_loader import (
TrafficItem,
TrafficSnapshot,
load_traffic_snapshot,
snapshot_to_dict,
)
from webui.traffic_views import render_traffic_page
from allocator_service import WorkCandidate
class TestTrafficClassification(unittest.TestCase):
def test_runnable_candidate_classification(self):
cand = WorkCandidate(
kind="issue",
number=640,
state="open",
labels=("status:ready",),
title="Web Console: Workflow traffic-control view (Phase 1)",
priority=20,
)
snap = load_traffic_snapshot(candidates=[cand])
self.assertEqual(len(snap.runnable), 1)
self.assertEqual(snap.runnable[0].number, 640)
self.assertTrue(snap.runnable[0].is_safe)
self.assertEqual(snap.runnable[0].traffic_state, "runnable")
def test_blocked_dependency_candidate_classification(self):
cand = WorkCandidate(
kind="issue",
number=643,
state="open",
labels=("status:ready",),
title="Web Console: Requests & intent preview (Phase 2)",
priority=20,
dependency_unmet=True,
dependency_reason="issue#643 depends on unresolved issue(s) #640; they are not closed",
)
snap = load_traffic_snapshot(candidates=[cand])
self.assertEqual(len(snap.blocked), 1)
self.assertEqual(snap.blocked[0].number, 643)
self.assertFalse(snap.blocked[0].is_safe)
self.assertEqual(snap.blocked[0].traffic_state, "blocked")
self.assertIn("depends on unresolved issue(s) #640", snap.blocked[0].block_reason)
def test_leased_candidate_classification(self):
cand = WorkCandidate(
kind="issue",
number=640,
state="open",
labels=("status:in-progress",),
title="Web Console: Workflow traffic-control view (Phase 1)",
priority=20,
)
lease = {
"kind": "issue",
"number": 640,
"session_id": "prgs-author-12345",
"role": "author",
"status": "active",
}
snap = load_traffic_snapshot(candidates=[cand], leases=[lease])
self.assertEqual(len(snap.leased), 1)
self.assertEqual(snap.leased[0].number, 640)
self.assertEqual(snap.leased[0].traffic_state, "leased")
self.assertIsNotNone(snap.leased[0].lease_info)
def test_needs_controller_candidate_classification(self):
cand = WorkCandidate(
kind="issue",
number=700,
state="open",
labels=("status:blocked",),
title="Controller intervention needed",
priority=10,
blocked=True,
)
snap = load_traffic_snapshot(candidates=[cand])
self.assertEqual(len(snap.needs_controller), 1)
self.assertEqual(snap.needs_controller[0].number, 700)
class TestTrafficLoader(unittest.TestCase):
def test_snapshot_to_dict_export(self):
cand = WorkCandidate(
kind="issue",
number=640,
state="open",
labels=("status:ready",),
title="Traffic control test",
priority=20,
)
snap = load_traffic_snapshot(candidates=[cand])
data = snapshot_to_dict(snap)
self.assertEqual(data["project_id"], "gitea-tools")
self.assertEqual(len(data["runnable"]), 1)
self.assertTrue(data["inventory_complete"])
def test_fail_closed_error_handling(self):
with mock.patch("webui.traffic_loader.load_queue_snapshot", side_effect=RuntimeError("Gitea connection failed")):
snap = load_traffic_snapshot()
self.assertIsNotNone(snap.fetch_error)
self.assertIn("Failed to load traffic state", snap.fetch_error)
self.assertEqual(len(snap.runnable), 0)
self.assertFalse(snap.inventory_complete)
class TestTrafficLivePath(unittest.TestCase):
"""Live path tests: inject QueueSnapshot + LeaseSnapshot (no candidates=).
Covers the production ``load_traffic_snapshot()`` branch that ``/traffic``
and ``/api/traffic`` actually execute (#640 B1B5).
"""
FULL_SHA = "069a9af7e6aa2c2994e07199d1b0814819457017"
def _queue(
self,
*,
prs=(),
issues=(),
):
from webui.queue_loader import QueueSnapshot
return QueueSnapshot(
project_id="gitea-tools",
repo_label="Scaled-Tech-Consulting/Gitea-Tools",
prs=tuple(prs),
issues=tuple(issues),
pr_pagination=None,
issue_pagination=None,
fetch_error=None,
)
def _lease(
self,
*,
claim_inventory=None,
reviewer_leases=(),
):
from webui.lease_loader import LeaseSnapshot
return LeaseSnapshot(
project_id="gitea-tools",
repo_label="Scaled-Tech-Consulting/Gitea-Tools",
issue_lock=None,
claim_inventory=claim_inventory or {"entries": [], "counts": {}},
reviewer_leases=tuple(reviewer_leases),
duplicate_prs=(),
duplicate_branches=(),
collision_history=(),
fetch_error=None,
)
def test_live_pr_uses_full_head_sha_and_is_runnable(self):
from webui.queue_loader import QueueItem
pr = QueueItem(
number=885,
title="traffic control",
badges=("in-review",),
extra={"head_sha": self.FULL_SHA[:12], "linked_issue": "640"},
signals={
"head_sha": self.FULL_SHA,
"mergeable": True,
"labels": (),
"linked_issue": 640,
},
)
q = self._queue(prs=[pr])
l = self._lease()
snap = load_traffic_snapshot(
fetch_queue_snapshot=lambda: q,
fetch_lease_snapshot=lambda: l,
)
self.assertIsNone(snap.fetch_error)
self.assertEqual(len(snap.runnable), 1)
item = snap.runnable[0]
self.assertEqual(item.kind, "pr")
self.assertEqual(item.number, 885)
self.assertEqual(item.head_sha, self.FULL_SHA)
self.assertNotEqual(item.head_sha, self.FULL_SHA[:12])
self.assertIsNone(item.block_reason)
self.assertEqual(len(snap.blocked), 0)
def test_live_pr_without_head_sha_is_blocked(self):
from webui.queue_loader import QueueItem
pr = QueueItem(
number=1,
title="missing pin",
badges=("open",),
extra={"head_sha": ""},
signals={"head_sha": "", "mergeable": True, "labels": ()},
)
snap = load_traffic_snapshot(
fetch_queue_snapshot=lambda: self._queue(prs=[pr]),
fetch_lease_snapshot=lambda: self._lease(),
)
self.assertEqual(len(snap.blocked) + len(snap.needs_controller), 1)
item = (snap.blocked or snap.needs_controller)[0]
self.assertIn("missing head_sha", (item.block_reason or "").lower())
def test_reviewer_lease_keys_by_pr_not_linked_issue(self):
from webui.queue_loader import QueueItem
pr = QueueItem(
number=885,
title="leased pr",
badges=("in-review",),
extra={"head_sha": self.FULL_SHA[:12]},
signals={"head_sha": self.FULL_SHA, "mergeable": True, "labels": ()},
)
issue = QueueItem(
number=640,
title="linked issue",
badges=("open",),
extra={},
signals={"labels": ()},
)
# Marker-shaped record: has both pr_number and issue_number; must
# attach to the PR only (B2).
reviewer_lease = {
"pr_number": 885,
"issue_number": 640,
"phase": "validating",
"reviewer_identity": "sysadmin",
"session_id": "review-sess-1",
}
snap = load_traffic_snapshot(
fetch_queue_snapshot=lambda: self._queue(prs=[pr], issues=[issue]),
fetch_lease_snapshot=lambda: self._lease(reviewer_leases=[reviewer_lease]),
)
leased_prs = [i for i in snap.leased if i.kind == "pr" and i.number == 885]
self.assertEqual(len(leased_prs), 1)
self.assertEqual(leased_prs[0].lease_info.get("pr_number"), 885)
# Issue 640 must not inherit the reviewer lease just because issue_number
# is present on the marker.
for item in list(snap.leased) + list(snap.runnable) + list(snap.blocked):
if item.kind == "issue" and item.number == 640:
self.assertIsNone(
item.lease_info,
"reviewer lease must not attach to linked issue #640",
)
break
else:
self.fail("expected issue #640 in traffic snapshot")
def test_claim_inventory_entries_key_marks_issue_leased(self):
from webui.queue_loader import QueueItem
issue = QueueItem(
number=640,
title="claimed issue",
badges=("claimed",),
extra={},
signals={"labels": ("status:in-progress",)},
)
inventory = {
"entries": [
{
"issue_number": 640,
"status": "active",
"latest_heartbeat": {"session_id": "author-sess-9"},
"reasons": ["claim has structured heartbeat proof"],
}
],
"counts": {"active": 1},
"in_progress_total": 1,
}
snap = load_traffic_snapshot(
fetch_queue_snapshot=lambda: self._queue(issues=[issue]),
fetch_lease_snapshot=lambda: self._lease(claim_inventory=inventory),
)
leased_issues = [i for i in snap.leased if i.kind == "issue" and i.number == 640]
self.assertEqual(len(leased_issues), 1)
self.assertEqual(leased_issues[0].traffic_state, "leased")
def test_active_claims_key_is_ignored(self):
"""B3 regression: fictional ``active_claims`` must not create lease_info."""
from webui.queue_loader import QueueItem
issue = QueueItem(
number=640,
title="open issue",
badges=("open",),
extra={},
signals={"labels": ()},
)
# Only the broken key — must NOT produce lease_info. Entries-less
# inventory is empty (entries is the real claim_inventory key).
inventory = {
"active_claims": [
{
"kind": "issue",
"number": 640,
"issue_number": 640,
"status": "active",
},
],
"counts": {},
}
snap = load_traffic_snapshot(
fetch_queue_snapshot=lambda: self._queue(issues=[issue]),
fetch_lease_snapshot=lambda: self._lease(claim_inventory=inventory),
)
items = [
i
for i in (
list(snap.runnable)
+ list(snap.leased)
+ list(snap.blocked)
+ list(snap.needs_controller)
)
if i.kind == "issue" and i.number == 640
]
self.assertEqual(len(items), 1)
self.assertIsNone(
items[0].lease_info,
"active_claims is not a real inventory key; entries-only",
)
class TestTrafficRoutesAndRendering(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_traffic_html_page_rendering(self):
cand1 = WorkCandidate(
kind="issue",
number=640,
state="open",
labels=("status:ready",),
title="Traffic View Implementation",
priority=20,
)
cand2 = WorkCandidate(
kind="issue",
number=643,
state="open",
labels=("status:ready",),
title="Dependent Feature",
priority=20,
dependency_unmet=True,
dependency_reason="issue#643 depends on unresolved issue(s) #640; they are not closed",
)
snap = load_traffic_snapshot(candidates=[cand1, cand2])
with mock.patch("webui.app.load_traffic_snapshot", return_value=snap):
response = self.client.get("/traffic")
self.assertEqual(response.status_code, 200)
self.assertIn("Workflow Traffic Control", response.text)
self.assertIn("1. Runnable Lanes", response.text)
self.assertIn("3. Blocked Items", response.text)
self.assertIn("Traffic View Implementation", response.text)
self.assertIn("depends on unresolved issue(s) #640", response.text)
def test_api_traffic_json_route(self):
cand = WorkCandidate(
kind="issue",
number=640,
state="open",
labels=("status:ready",),
title="Traffic View API Test",
priority=20,
)
snap = load_traffic_snapshot(candidates=[cand])
with mock.patch("webui.app.load_traffic_snapshot", return_value=snap):
response = self.client.get("/api/traffic")
self.assertEqual(response.status_code, 200)
data = response.json()
self.assertEqual(data["project_id"], "gitea-tools")
self.assertEqual(len(data["runnable"]), 1)
self.assertEqual(data["runnable"][0]["number"], 640)
def test_render_traffic_fail_closed_page(self):
snap = TrafficSnapshot(
project_id="gitea-tools",
repo_label="Scaled-Tech-Consulting/Gitea-Tools",
runnable=(),
leased=(),
blocked=(),
needs_controller=(),
terminal_complete=(),
next_roles=(),
fetch_error="Gitea credentials unavailable for gitea.prgs.cc",
inventory_complete=False,
)
html = render_traffic_page(snap)
self.assertIn("Traffic data unavailable", html)
self.assertIn("Fail closed", html)
self.assertNotIn("1. Runnable Lanes", html)
if __name__ == "__main__":
unittest.main()
+14
View File
@@ -42,6 +42,8 @@ from webui.lease_loader import load_lease_snapshot, snapshot_to_dict as lease_sn
from webui.lease_views import render_leases_page from webui.lease_views import render_leases_page
from webui.queue_loader import load_queue_snapshot, snapshot_to_dict as queue_snapshot_to_dict from webui.queue_loader import load_queue_snapshot, snapshot_to_dict as queue_snapshot_to_dict
from webui.queue_views import render_queue_page from webui.queue_views import render_queue_page
from webui.traffic_loader import load_traffic_snapshot, snapshot_to_dict as traffic_snapshot_to_dict
from webui.traffic_views import render_traffic_page
from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as worktree_snapshot_to_dict from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as worktree_snapshot_to_dict
from webui.worktree_views import render_worktrees_page from webui.worktree_views import render_worktrees_page
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
@@ -80,6 +82,7 @@ def _stub_page(title: str, description: str) -> HTMLResponse:
_LEGACY_PAGES = ( _LEGACY_PAGES = (
("/traffic", "Traffic", "workflow traffic-control view (#640)"),
("/queue", "Queue", "live PR and issue dashboard (#429)"), ("/queue", "Queue", "live PR and issue dashboard (#429)"),
("/projects", "Projects", "registry and onboarding (#427)"), ("/projects", "Projects", "registry and onboarding (#427)"),
("/prompts", "Prompts", "canonical workflow prompt library (#428)"), ("/prompts", "Prompts", "canonical workflow prompt library (#428)"),
@@ -200,6 +203,15 @@ async def api_queue(_request: Request) -> JSONResponse:
return JSONResponse(queue_snapshot_to_dict(load_queue_snapshot())) return JSONResponse(queue_snapshot_to_dict(load_queue_snapshot()))
async def traffic(_request: Request) -> HTMLResponse:
snapshot = load_traffic_snapshot()
return HTMLResponse(render_traffic_page(snapshot))
async def api_traffic(_request: Request) -> JSONResponse:
return JSONResponse(traffic_snapshot_to_dict(load_traffic_snapshot()))
def _load_project_registry() -> tuple[ProjectRegistry | None, RegistryError | None]: def _load_project_registry() -> tuple[ProjectRegistry | None, RegistryError | None]:
"""Load the registry, converting validation failure into a fail-closed pair.""" """Load the registry, converting validation failure into a fail-closed pair."""
try: try:
@@ -736,6 +748,8 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/system-health", system_health, methods=["GET"]), Route("/system-health", system_health, methods=["GET"]),
Route("/queue", queue, methods=["GET"]), Route("/queue", queue, methods=["GET"]),
Route("/api/queue", api_queue, methods=["GET"]), Route("/api/queue", api_queue, methods=["GET"]),
Route("/traffic", traffic, methods=["GET"]),
Route("/api/traffic", api_traffic, methods=["GET"]),
Route("/projects", projects, methods=["GET"]), Route("/projects", projects, methods=["GET"]),
Route("/projects/{project_id}", project_detail, methods=["GET"]), Route("/projects/{project_id}", project_detail, methods=["GET"]),
Route("/api/projects", api_projects, methods=["GET"]), Route("/api/projects", api_projects, methods=["GET"]),
+8 -1
View File
@@ -182,10 +182,17 @@ def _extract_reviewer_leases(
parsed = parse_reviewer_lease_comment(comment.get("body") or "") parsed = parse_reviewer_lease_comment(comment.get("body") or "")
if not parsed: if not parsed:
continue continue
subject_pr = parsed.get("pr_number") or pr_number
leases.append( leases.append(
{ {
**parsed, **parsed,
"pr_number": parsed.get("pr_number") or pr_number, "pr_number": subject_pr,
# The lease subject is the PR, never the linked issue: a
# reviewer lease on PR #N must not be attributed to issue #N
# or to the issue that PR closes (#640).
"kind": "pr",
"number": subject_pr,
"role": "reviewer",
"comment_id": comment.get("id"), "comment_id": comment.get("id"),
"author": (comment.get("user") or {}).get("login"), "author": (comment.get("user") or {}).get("login"),
"created_at": comment.get("created_at"), "created_at": comment.get("created_at"),
+1
View File
@@ -41,6 +41,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
NavItem("/system-health", "System health"), NavItem("/system-health", "System health"),
)), )),
NavGroup("Traffic", ( NavGroup("Traffic", (
NavItem("/traffic", "Traffic control"),
NavItem("/queue", "Queue"), NavItem("/queue", "Queue"),
NavItem("/leases", "Leases"), NavItem("/leases", "Leases"),
NavItem("/actions", "Actions"), NavItem("/actions", "Actions"),
+34 -3
View File
@@ -4,7 +4,7 @@ from __future__ import annotations
import os import os
import re import re
from dataclasses import dataclass from dataclasses import dataclass, field
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Any, Callable from typing import Any, Callable
from urllib.parse import urlparse from urllib.parse import urlparse
@@ -31,10 +31,20 @@ class PaginationMeta:
@dataclass(frozen=True) @dataclass(frozen=True)
class QueueItem: class QueueItem:
"""One queue row.
``extra`` holds *display* strings for the queue page (values are truncated
or humanized for rendering). ``signals`` holds the *authoritative* typed
values taken straight from the Gitea payload, for consumers that classify
or pin state rather than render it (#640). Never derive identity or
concurrency decisions from ``extra``.
"""
number: int number: int
title: str title: str
badges: tuple[str, ...] badges: tuple[str, ...]
extra: dict[str, str] extra: dict[str, str]
signals: dict[str, Any] = field(default_factory=dict)
@dataclass(frozen=True) @dataclass(frozen=True)
@@ -134,21 +144,37 @@ def _format_pr_item(pr: dict, badges: tuple[str, ...]) -> QueueItem:
"mergeable" if mergeable is True else "conflicted" if mergeable is False else "unknown" "mergeable" if mergeable is True else "conflicted" if mergeable is False else "unknown"
) )
linked = _extract_linked_issue(pr.get("title"), pr.get("body")) linked = _extract_linked_issue(pr.get("title"), pr.get("body"))
head_sha = str(head.get("sha") or "")
labels = tuple(
str(lb.get("name") or "") for lb in (pr.get("labels") or []) if lb.get("name")
)
return QueueItem( return QueueItem(
number=int(pr["number"]), number=int(pr["number"]),
title=str(pr.get("title") or ""), title=str(pr.get("title") or ""),
badges=badges, badges=badges,
extra={ extra={
"branch": f"{head.get('ref', '?')}{base.get('ref', '?')}", "branch": f"{head.get('ref', '?')}{base.get('ref', '?')}",
"head_sha": str(head.get("sha") or "")[:12], # Display only — truncated. Pin against signals["head_sha"] instead.
"head_sha": head_sha[:12],
"mergeable": merge_label, "mergeable": merge_label,
"linked_issue": str(linked) if linked is not None else "", "linked_issue": str(linked) if linked is not None else "",
}, },
signals={
"head_sha": head_sha,
"head_ref": str(head.get("ref") or ""),
"base_ref": str(base.get("ref") or ""),
"mergeable": mergeable if isinstance(mergeable, bool) else None,
"labels": labels,
"linked_issue": linked,
},
) )
def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem: def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
labels = ", ".join(lb.get("name", "") for lb in issue.get("labels", [])) label_names = tuple(
str(lb.get("name") or "") for lb in (issue.get("labels") or []) if lb.get("name")
)
labels = ", ".join(label_names)
assignee = (issue.get("assignee") or {}).get("login", "") assignee = (issue.get("assignee") or {}).get("login", "")
return QueueItem( return QueueItem(
number=int(issue["number"]), number=int(issue["number"]),
@@ -159,6 +185,11 @@ def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
"assignee": assignee or "unassigned", "assignee": assignee or "unassigned",
"state": str(issue.get("state") or ""), "state": str(issue.get("state") or ""),
}, },
signals={
"labels": label_names,
"assignee": assignee,
"state": str(issue.get("state") or ""),
},
) )
+449
View File
@@ -0,0 +1,449 @@
"""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()
+170
View File
@@ -0,0 +1,170 @@
"""HTML rendering for Phase 1 Traffic-Control View (#640)."""
from __future__ import annotations
from html import escape
from typing import Sequence
from webui.layout import render_page
from webui.traffic_loader import TrafficItem, TrafficSnapshot
def _render_badges(badges: Sequence[str]) -> str:
if not badges:
return ""
out = []
for b in badges:
cls = "badge"
b_lower = b.lower()
if "blocked" in b_lower or "unmet" in b_lower:
cls += " badge-blocked"
elif "claimed" in b_lower or "in-progress" in b_lower or "leased" in b_lower:
cls += " badge-claimed"
elif "review" in b_lower or "ready" in b_lower:
cls += " badge-in-review"
elif "duplicate" in b_lower:
cls += " badge-duplicate"
elif "stale" in b_lower:
cls += " badge-stale"
out.append(f'<span class="{cls}">{escape(b)}</span>')
return f'<div class="badges">{"".join(out)}</div>'
def _render_traffic_item_row(item: TrafficItem) -> str:
kind_label = escape(item.kind.upper())
num_str = f"#{item.number}"
title_str = escape(item.title)
role_str = escape(item.expected_role)
badges_html = _render_badges(item.badges)
reason_html = ""
if item.block_reason:
reason_html = f'<div class="muted" style="font-size:0.82rem; margin-top:0.2rem;"><strong>Blocker:</strong> {escape(item.block_reason)}</div>'
lease_html = ""
if item.lease_info:
owner = escape(str(item.lease_info.get("session_id") or item.lease_info.get("reviewer_identity") or "active worker"))
lease_html = f'<div class="muted" style="font-size:0.82rem; margin-top:0.2rem;"><strong>Lease:</strong> {owner}</div>'
return f"""<tr>
<td><code>{kind_label} {num_str}</code></td>
<td>
<div><strong>{title_str}</strong> {badges_html}</div>
{reason_html}
{lease_html}
</td>
<td><code>{role_str}</code></td>
</tr>"""
def _render_traffic_table(items: Sequence[TrafficItem], empty_message: str) -> str:
if not items:
return f'<p class="muted">{escape(empty_message)}</p>'
rows = "".join(_render_traffic_item_row(item) for item in items)
return f"""<table class="registry">
<thead>
<tr>
<th style="width: 15%;">Item</th>
<th style="width: 65%;">Title & Details</th>
<th style="width: 20%;">Next Role</th>
</tr>
</thead>
<tbody>
{rows}
</tbody>
</table>"""
def _render_next_roles(next_roles: Sequence[dict]) -> str:
if not next_roles:
return ""
cards = []
for r in next_roles:
role = escape(r.get("role", "unknown"))
status = r.get("status", "idle")
prompt = escape(r.get("prompt", ""))
status_cls = "badge-health-ok" if status == "safe" else ("badge-blocked" if "blocked" in status else "badge-health-skipped")
cards.append(f"""<div class="health-card" style="margin-bottom:0.75rem;">
<div style="display:flex; justify-content:space-between; align-items:center;">
<h3>Role: <code>{role}</code></h3>
<span class="badge {status_cls}">status: {escape(status)}</span>
</div>
<p class="meta" style="margin:0.35rem 0 0;">{prompt}</p>
</div>""")
return f"""<div style="margin: 1.5rem 0;">
<h3>Next Safe Role Actions</h3>
{"".join(cards)}
</div>"""
def render_traffic_page(snapshot: TrafficSnapshot) -> str:
"""Render the full HTML view for workflow traffic control."""
if snapshot.fetch_error:
body = f"""<h2>Workflow Traffic Control</h2>
<p class="meta">Repository: <code>{escape(snapshot.repo_label)}</code></p>
<div class="health-card health-stale">
<h3>Traffic data unavailable</h3>
<p class="health-headline">{escape(snapshot.fetch_error)}</p>
<p class="muted">Fail closed: traffic state cannot be established cleanly. Check credentials or remote connectivity.</p>
</div>"""
return render_page(title="Traffic Control", body_html=body)
runnable_count = len(snapshot.runnable)
leased_count = len(snapshot.leased)
blocked_count = len(snapshot.blocked)
controller_count = len(snapshot.needs_controller)
terminal_count = len(snapshot.terminal_complete)
summary_bar = f"""<div class="health-card" style="display:flex; flex-wrap:wrap; gap:1rem; align-items:center;">
<div><strong>Runnable:</strong> <span class="badge badge-health-ok">{runnable_count}</span></div>
<div><strong>Leased:</strong> <span class="badge badge-claimed">{leased_count}</span></div>
<div><strong>Blocked:</strong> <span class="badge badge-blocked">{blocked_count}</span></div>
<div><strong>Needs Controller:</strong> <span class="badge badge-duplicate">{controller_count}</span></div>
<div><strong>Terminal Complete:</strong> <span class="badge badge-stale">{terminal_count}</span></div>
</div>"""
next_roles_html = _render_next_roles(snapshot.next_roles)
sections_html = f"""
<div class="prompt-card">
<h3>1. Runnable Lanes (Ready for Allocation)</h3>
<p class="muted">Safe work items with no unmet dependencies or active leases. Safe for allocation.</p>
{_render_traffic_table(snapshot.runnable, "No runnable items ready for allocation.")}
</div>
<div class="prompt-card">
<h3>2. In-Progress Work (Active Leases)</h3>
<p class="muted">Work items currently leased and actively being worked by an assigned role session.</p>
{_render_traffic_table(snapshot.leased, "No active leases in flight.")}
</div>
<div class="prompt-card">
<h3>3. Blocked Items (Dependencies / Locks)</h3>
<p class="muted">Items blocked by unmet dependency issues, a missing head pin, a merge conflict, or an active terminal review lock. Items labelled status:blocked route to section 4. Never presented as safe.</p>
{_render_traffic_table(snapshot.blocked, "No blocked items.")}
</div>
<div class="prompt-card">
<h3>4. Needs Controller Intervention</h3>
<p class="muted">Items requiring controller routing, diagnosis, or cross-role assignment.</p>
{_render_traffic_table(snapshot.needs_controller, "No items requiring controller intervention.")}
</div>
<div class="prompt-card">
<h3>5. Terminal / Complete Candidates</h3>
<p class="muted">Items ready for terminal reconciliation or post-merge worktree cleanup.</p>
{_render_traffic_table(snapshot.terminal_complete, "No terminal complete candidates.")}
</div>
"""
body = f"""<h2>Workflow Traffic Control</h2>
<p class="meta">Repository: <code>{escape(snapshot.repo_label)}</code></p>
{summary_bar}
{next_roles_html}
{sections_html}"""
return render_page(title="Traffic Control", body_html=body)