Author SHA1 Message Date
sysadminandClaude Opus 4.8 1948d3dc21 fix(webui): repair traffic live path contracts for #640 review
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]>
2026-07-24 18:31:28 -04:00
sysadmin 069a9af7e6 feat(webui): implement workflow traffic-control view (Closes #640) 2026-07-24 17:34:50 -04:00
10 changed files with 1127 additions and 463 deletions
+17 -81
View File
@@ -133,15 +133,10 @@ ROLE_ACTIONS: dict[str, tuple[tuple[str, ...], tuple[str, ...]]] = {
# Body phrases that prove an issue is an implementation container, not a
# unit of direct author work (#844 / #854). Matched case-insensitively against
# the issue body. Title alone is never sufficient (ordinary issues may mention
# "epic", "roadmap", "vision", or "umbrella" incidentally).
#
# #854 extends the #844 marker set so product-vision (#652), phased-roadmap
# (#653), and umbrella (#655) coordination records — which do not use the word
# "epic" — are classified with the same semantic exclusion as epic containers.
# unit of direct author work (#844). Matched case-insensitively against the
# issue body. Title alone is never sufficient (ordinary issues may mention
# "epic" incidentally).
_CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
# Epic / child-only (#844, live #631)
"implementation is delivered via child issues only",
"implementation is delivered through child issues only",
"implementation is delivered via child issues",
@@ -155,26 +150,9 @@ _CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
"coordination container",
"child-only container",
"implementation is delegated to child",
# Vision / roadmap / umbrella coordination (#854, live #652/#653/#655).
# Prefer authoritative non-implementation / child-only scope language over
# bare words like "roadmap" so ordinary implementable issues that mention
# a parent vision or roadmap stay eligible.
"do not implement features on this issue",
"implementing features on this roadmap issue",
"implementation is via linked children only",
"no product feature claimed complete on this issue alone",
"phased delivery roadmap and epic sequencing",
"this issue is the enduring source of truth",
"enduring source of truth for the",
"canonical product vision — enduring source of truth",
"canonical product vision - enduring source of truth",
"state: vision-active",
"state: roadmap-active",
)
# Explicit epic / umbrella / vision / roadmap labels (structured evidence
# preferred over title). Tracker alone is *not* included — ordinary issues
# may carry a tracker label without being non-implementable containers.
# Explicit epic / umbrella labels (structured evidence preferred over title).
_EPIC_LABELS: frozenset[str] = frozenset(
{
"type:epic",
@@ -183,16 +161,6 @@ _EPIC_LABELS: frozenset[str] = frozenset(
"scope:epic",
"type:umbrella",
"umbrella",
"kind:umbrella",
"scope:umbrella",
"type:vision",
"vision",
"kind:vision",
"scope:vision",
"type:roadmap",
"roadmap",
"kind:roadmap",
"scope:roadmap",
}
)
@@ -254,45 +222,16 @@ class WorkCandidate:
}
def _title_container_prefix(title_l: str) -> str | None:
"""Return a coordination-title prefix token if *title_l* uses one (#854).
Title prefixes alone never exclude; they only corroborate body/label
evidence. Ordinary issues may say "roadmap" or "vision" mid-title.
"""
for prefix, token in (
("epic:", "title_epic_prefix"),
("epic ", "title_epic_prefix"),
("umbrella:", "title_umbrella_prefix"),
("umbrella ", "title_umbrella_prefix"),
("roadmap:", "title_roadmap_prefix"),
("roadmap ", "title_roadmap_prefix"),
("product vision:", "title_vision_prefix"),
("product vision ", "title_vision_prefix"),
("vision:", "title_vision_prefix"),
("vision ", "title_vision_prefix"),
):
if title_l.startswith(prefix):
return token
return None
def classify_epic_or_child_only_container(
c: WorkCandidate,
) -> tuple[bool, str | None]:
"""Return whether *c* is a non-implementable coordination container (#844/#854).
"""Return whether *c* is an epic / child-only implementation container (#844).
Exclusion uses structured evidence first (labels, body scope language).
A bare title containing the words "epic", "roadmap", "vision", or
"umbrella" is **not** enough — ordinary implementable issues may mention
those terms incidentally. Explicit title prefixes (``Epic:``, ``Roadmap:``,
``Product vision:``, ``Umbrella:``) only count when the body also proves
child-only / no-direct-implementation scope (or a container label is
present).
Covers epic, product-vision, phased-roadmap, umbrella, and child-only
records so the allocator never assigns coordination containers as direct
author work.
A bare title containing the word "epic" is **not** enough — ordinary
implementable issues may mention epics incidentally. A title that is
explicitly prefixed ``Epic:`` only counts when the body also proves
child-only / no-direct-implementation scope (or an epic label is present).
PRs are never classified as containers here (they already have a head).
"""
@@ -306,28 +245,25 @@ def classify_epic_or_child_only_container(
title_l = title.lower()
body_hits = [m for m in _CHILD_ONLY_BODY_MARKERS if m in body_l]
title_prefix = _title_container_prefix(title_l)
title_epic_prefix = title_l.startswith("epic:") or title_l.startswith("epic ")
if epic_label:
detail = f"label={epic_label[0]}"
if body_hits:
detail = f"{detail}; body_marker={body_hits[0]!r}"
if title_prefix:
detail = f"{title_prefix}; {detail}"
return True, detail
if body_hits:
# Body proves child-only / vision / roadmap / umbrella scope. Title
# prefixes are corroborating but not required — containers without the
# title word still exclude.
# Body proves child-only / umbrella scope. Title "Epic:" is corroborating
# but not required — containers without the word still exclude.
detail = f"body_marker={body_hits[0]!r}"
if title_prefix:
detail = f"{title_prefix}; {detail}"
if title_epic_prefix:
detail = f"title_epic_prefix; {detail}"
return True, detail
# Title-only coordination prefix without body scope evidence is
# insufficient (#844/#854 AC: eligibility does not rely solely on a title
# word). Incidental mid-title mentions without markers stay eligible.
# Title-only "Epic:" without body scope evidence is insufficient (#844 AC:
# eligibility does not rely solely on the word "Epic" in a title).
# Similarly, incidental "epic" mid-title without markers stays eligible.
return False, None
+27
View File
@@ -57,6 +57,33 @@ 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) |
| `/queue` | Live PR and issue queue dashboard (#429) |
| `/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 |
### 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 status:blocked | Author/controller remediation first |
| **needs_controller** | Contaminated / controller-only diagnosis | Controller only |
| **terminal_complete** | Reconciler / terminal-lock territory | Reconciler cleanup path |
**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.
| `/projects` | Project registry list with status and onboarding progress (#427, #635) |
| `/projects/{id}` | Project detail + onboarding checklist |
| `/api/v1/projects` | Versioned JSON registry export (#635) |
@@ -1,378 +0,0 @@
"""Allocator semantic container exclusion for vision/roadmap/umbrella (#854).
#844 excluded epic / child-only containers (live #631) but product-vision
(#652), phased-roadmap (#653), and umbrella (#655) coordination records still
ranked as implementable work. This module is the live-equivalent canary:
* #631 / #652 / #653 / #655-shaped records are all excluded in one inventory.
* Independently executable children remain eligible and can be selected.
* Ordinary issues that merely mention vision / roadmap / umbrella stay eligible.
* Excluded containers never receive assignments or workflow leases.
* Structured skip reason ``epic_or_child_only_container`` is reported.
* Candidate-set fingerprint remains stable after exclusions.
"""
from __future__ import annotations
import os
import tempfile
import unittest
from allocator_service import (
OUTCOME_ASSIGNED,
OUTCOME_PREVIEW,
SKIP_EPIC_OR_CHILD_ONLY_CONTAINER,
WorkCandidate,
allocate_next_work,
candidate_set_fingerprint,
classify_epic_or_child_only_container,
)
from control_plane_db import ControlPlaneDB
REMOTE = "prgs"
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
# Minimal bodies mirroring live coordination records (not full issue text).
_EPIC_631_BODY = """
## Scope (umbrella)
This epic owns the **product roadmap and linkage** for the Web Console.
Implementation is delivered via child issues only.
## Explicit non-goals
* Do not implement product features in this epic issue itself.
* No product feature implementation is claimed complete solely on this epic.
"""
_VISION_652_BODY = """
## Canonical product vision — enduring source of truth
**This issue is the enduring source of truth for the MCP Control Plane Web Console product vision.**
## Implementation linkage
* **Do not implement features on this issue.**
* Sequencing: roadmap issue + #631 children.
## Canonical issue state
```text
STATE: vision-active
WHO_IS_NEXT: controller (triage/ordering) / author (implementation of linked children only)
```
"""
_ROADMAP_653_BODY = """
## Purpose
This issue is the **phased delivery roadmap and epic sequencing** for the MCP Control Plane Web Console.
## Non-goals
* Implementing features on this roadmap issue.
* Deleting vision items by omitting them from phases without #652 change log.
## Canonical issue state
```text
STATE: roadmap-active
WHO_IS_NEXT: author
```
"""
_UMBRELLA_655_BODY = """
## Scope (umbrella)
This issue owns the **canonical restart-governance program**. Implementation is via linked children only.
## Acceptance criteria (umbrella)
6. No product feature claimed complete on this issue alone.
"""
_CHILD_BODY = """
## Problem
Operators need a workflow-event timeline model for Phase 1.
## Acceptance criteria
- [ ] Timeline model API exists
"""
def _issue(
number: int,
*,
title: str = "",
body: str = "",
labels: tuple[str, ...] = ("status:ready", "type:feature"),
priority: int = 20,
) -> WorkCandidate:
return WorkCandidate(
kind="issue",
number=number,
state="open",
labels=labels,
title=title or f"issue {number}",
body=body,
priority=priority,
)
def _live_shaped_containers() -> list[WorkCandidate]:
return [
_issue(
631,
title="Epic: MCP Control Plane Web Console",
body=_EPIC_631_BODY,
),
_issue(
652,
title="Product vision: MCP Control Plane Web Console (canonical)",
body=_VISION_652_BODY,
),
_issue(
653,
title="Roadmap: MCP Control Plane Web Console (phased delivery)",
body=_ROADMAP_653_BODY,
),
_issue(
655,
title="Umbrella: Governed MCP restart coordination and zero-disruption recovery",
body=_UMBRELLA_655_BODY,
),
]
class ClassifySemanticContainersTest(unittest.TestCase):
def test_652_vision_is_container(self) -> None:
c = _issue(
652,
title="Product vision: MCP Control Plane Web Console (canonical)",
body=_VISION_652_BODY,
)
is_c, detail = classify_epic_or_child_only_container(c)
self.assertTrue(is_c)
self.assertIsNotNone(detail)
self.assertIn("body_marker", detail or "")
def test_653_roadmap_is_container(self) -> None:
c = _issue(
653,
title="Roadmap: MCP Control Plane Web Console (phased delivery)",
body=_ROADMAP_653_BODY,
)
is_c, detail = classify_epic_or_child_only_container(c)
self.assertTrue(is_c)
self.assertIn("body_marker", detail or "")
def test_655_umbrella_is_container(self) -> None:
c = _issue(
655,
title="Umbrella: Governed MCP restart coordination and zero-disruption recovery",
body=_UMBRELLA_655_BODY,
)
is_c, detail = classify_epic_or_child_only_container(c)
self.assertTrue(is_c)
self.assertIn("body_marker", detail or "")
def test_631_still_container_after_854(self) -> None:
c = _issue(
631,
title="Epic: MCP Control Plane Web Console",
body=_EPIC_631_BODY,
)
is_c, detail = classify_epic_or_child_only_container(c)
self.assertTrue(is_c)
self.assertIn("body_marker", detail or "")
def test_incidental_vision_roadmap_umbrella_words_not_container(self) -> None:
cases = (
(
"Document vision handoff conventions",
"Update docs so implementable issues that mention a vision "
"remain independently executable.",
),
(
"Clarify roadmap sequencing notes",
"Write a short note about how the roadmap issue relates to children.",
),
(
"Umbrella recovery checklist for authors",
"Authors should still implement the concrete recovery fix here.",
),
(
"Product vision wording in the help text",
"Fix a typo in the operator-facing help string that says product vision.",
),
)
for title, body in cases:
with self.subTest(title=title):
c = _issue(900, title=title, body=body)
is_c, detail = classify_epic_or_child_only_container(c)
self.assertFalse(is_c)
self.assertIsNone(detail)
def test_title_prefix_alone_not_container(self) -> None:
for title in (
"Epic: something mentioned only in title",
"Roadmap: title only without body scope",
"Product vision: title only without body scope",
"Umbrella: title only without body scope",
):
with self.subTest(title=title):
c = _issue(
901,
title=title,
body="Implement a concrete fix for the allocator skip list.",
)
is_c, detail = classify_epic_or_child_only_container(c)
self.assertFalse(is_c)
self.assertIsNone(detail)
def test_roadmap_label_alone_is_container(self) -> None:
c = _issue(
902,
title="Console delivery sequencing",
body="Track phased delivery only.",
labels=("status:ready", "type:roadmap"),
)
is_c, detail = classify_epic_or_child_only_container(c)
self.assertTrue(is_c)
self.assertIn("type:roadmap", detail or "")
def test_child_referencing_parent_policy_stays_eligible(self) -> None:
"""Children may quote parent policy without becoming containers."""
c = _issue(
637,
title="Web Console: Workflow-event timeline model (Phase 1)",
body=(
_CHILD_BODY
+ "\n\nParent #652 says do not implement on the vision issue; "
"this child is the implementable unit."
),
)
is_c, _ = classify_epic_or_child_only_container(c)
self.assertFalse(is_c)
class AllocateSemanticContainerExclusionTest(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
def _alloc(self, candidates, **kwargs):
defaults = dict(
session_id="sess-854",
role="author",
remote=REMOTE,
org=ORG,
repo=REPO,
profile_name="prgs-author",
username="jcwalker3",
claims={},
apply=False,
)
defaults.update(kwargs)
return allocate_next_work(self.db, candidates=candidates, **defaults)
def test_live_equivalent_canary_excludes_all_containers_selects_child(self) -> None:
containers = _live_shaped_containers()
child = _issue(
637,
title="Web Console: Workflow-event timeline model (Phase 1)",
body=_CHILD_BODY,
)
inventory = containers + [child]
res = self._alloc(inventory, apply=False)
self.assertTrue(res["success"], res)
self.assertEqual(res["outcome"], OUTCOME_PREVIEW)
self.assertEqual(res["selected"]["number"], 637)
skipped = {s["number"]: s for s in res["skipped"]}
for number in (631, 652, 653, 655):
self.assertIn(number, skipped, res["skipped"])
self.assertEqual(
skipped[number]["reason_code"],
SKIP_EPIC_OR_CHILD_ONLY_CONTAINER,
)
self.assertIn(
SKIP_EPIC_OR_CHILD_ONLY_CONTAINER, skipped[number]["reason"]
)
def test_containers_cannot_receive_assignment_or_lease(self) -> None:
containers = _live_shaped_containers()
res = self._alloc(containers, apply=True)
self.assertTrue(res["success"], res)
self.assertNotEqual(res["outcome"], OUTCOME_ASSIGNED)
self.assertIsNone(res.get("assignment"))
self.assertIsNone(res.get("selected"))
skipped = {s["number"]: s for s in res["skipped"]}
for number in (631, 652, 653, 655):
self.assertEqual(
skipped[number]["reason_code"],
SKIP_EPIC_OR_CHILD_ONLY_CONTAINER,
)
leases = []
if hasattr(self.db, "list_active_leases"):
leases = self.db.list_active_leases(
remote=REMOTE, org=ORG, repo=REPO
)
if not leases and hasattr(self.db, "list_leases"):
leases = self.db.list_leases(remote=REMOTE, org=ORG, repo=REPO)
for lease in leases or []:
work_number = (
lease.get("work_number") if isinstance(lease, dict) else None
)
self.assertNotIn(work_number, {631, 652, 653, 655})
def test_apply_selects_child_not_container(self) -> None:
containers = _live_shaped_containers()
child = _issue(
637,
title="Web Console: Workflow-event timeline model (Phase 1)",
body=_CHILD_BODY,
)
res = self._alloc(containers + [child], apply=True)
self.assertTrue(res["success"], res)
self.assertEqual(res["outcome"], OUTCOME_ASSIGNED)
self.assertEqual(res["selected"]["number"], 637)
self.assertEqual(res["assignment"]["work_number"], 637)
def test_fingerprint_stable_with_containers_present(self) -> None:
containers = _live_shaped_containers()
child = _issue(
637,
title="Web Console: Workflow-event timeline model (Phase 1)",
body=_CHILD_BODY,
)
inventory = containers + [child]
fp_before = candidate_set_fingerprint(inventory)
res = self._alloc(inventory, apply=False)
self.assertTrue(res["success"], res)
self.assertEqual(res["selected"]["number"], 637)
# Allocator reports the same CAS fingerprint for the full candidate set.
reported = res.get("candidate_set_fingerprint")
self.assertEqual(reported, fp_before)
# Re-fingerprint of the same inventory is byte-stable.
self.assertEqual(candidate_set_fingerprint(inventory), fp_before)
def test_incidental_mentions_remain_eligible(self) -> None:
ordinary = _issue(
700,
title="Document roadmap handoff conventions",
body="Write runbook text about vision vs roadmap vs child issues.",
)
res = self._alloc([ordinary], apply=False)
self.assertTrue(res["success"], res)
self.assertEqual(res["selected"]["number"], 700)
self.assertEqual(res["skipped"], [])
if __name__ == "__main__":
unittest.main()
+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()
+13
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.queue_loader import load_queue_snapshot, snapshot_to_dict as queue_snapshot_to_dict
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_views import render_worktrees_page
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
@@ -200,6 +202,15 @@ async def api_queue(_request: Request) -> JSONResponse:
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]:
"""Load the registry, converting validation failure into a fail-closed pair."""
try:
@@ -736,6 +747,8 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/system-health", system_health, methods=["GET"]),
Route("/queue", 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/{project_id}", project_detail, 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 "")
if not parsed:
continue
subject_pr = parsed.get("pr_number") or pr_number
leases.append(
{
**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"),
"author": (comment.get("user") or {}).get("login"),
"created_at": comment.get("created_at"),
+1
View File
@@ -41,6 +41,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
NavItem("/system-health", "System health"),
)),
NavGroup("Traffic", (
NavItem("/traffic", "Traffic control"),
NavItem("/queue", "Queue"),
NavItem("/leases", "Leases"),
NavItem("/actions", "Actions"),
+34 -3
View File
@@ -4,7 +4,7 @@ from __future__ import annotations
import os
import re
from dataclasses import dataclass
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Callable
from urllib.parse import urlparse
@@ -31,10 +31,20 @@ class PaginationMeta:
@dataclass(frozen=True)
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
title: str
badges: tuple[str, ...]
extra: dict[str, str]
signals: dict[str, Any] = field(default_factory=dict)
@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"
)
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(
number=int(pr["number"]),
title=str(pr.get("title") or ""),
badges=badges,
extra={
"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,
"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:
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", "")
return QueueItem(
number=int(issue["number"]),
@@ -159,6 +185,11 @@ def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
"assignee": assignee or "unassigned",
"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, status:blocked, or active terminal review locks. 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)