Author SHA1 Message Date
jcwalker3 a81db75402 Merge branch 'master' into feat/issue-641-runtime-session-view 2026-07-24 22:34:40 -05:00
sysadmin 7af40fb5ff Merge pull request 'fix(allocator): exclude vision/roadmap/umbrella coordination containers (Closes #854)' (#883) from fix/issue-854-semantic-container-exclusion into master 2026-07-24 22:27:58 -05:00
sysadminandClaude Opus 4.8 619f679077 feat(webui): Runtime and session view (Phase 1) (Closes #641)
Compose runtime health with inventory sessions/namespaces/worktrees into a
live /sessions page and JSON API. Surface stale PID/lease flags and durable
contamination markers when detectable. Recovery links name sanctioned
reconnect/restart paths only — no kill controls.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-24 22:45:27 -04:00
jcwalker3 3b68d15593 Merge branch 'master' into fix/issue-854-semantic-container-exclusion 2026-07-24 21:27:55 -05:00
sysadmin a4c73766f4 Merge pull request 'feat(webui): Workflow traffic-control view (Phase 1) (Closes #640)' (#885) from issue-640 into master 2026-07-24 21:10:51 -05:00
jcwalker3 b2e28428a4 Merge branch 'master' into fix/issue-854-semantic-container-exclusion 2026-07-24 21:06:41 -05:00
jcwalker3 95e4aae287 Merge branch 'master' into issue-640 2026-07-24 21:06:33 -05:00
sysadmin dac40ab9b3 docs(webui): update traffic state vocabulary docs and app nav for #640 2026-07-24 21:33:29 -04:00
sysadmin ccde9e8f11 Merge pull request 'fix(webui): migrate Starlette TestClient to httpx2 (Closes #682)' (#884) from fix/issue-682-starlette-httpx2 into master 2026-07-24 18:11:02 -05:00
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
sysadmin 0a78da39e5 fix(allocator): exclude vision/roadmap/umbrella coordination containers (#854)
#844 only caught epic-shaped child-only records. Live allocation still
selected product vision (#652), phased roadmap (#653), and umbrella (#655)
as implement targets. Extend pre-rank semantic classification with body
markers and container labels for those coordination records, keep title-
only and incidental mentions eligible, and add a live-equivalent canary.

Closes #854
2026-07-24 17:17:05 -04:00
15 changed files with 2854 additions and 36 deletions
+81 -17
View File
@@ -133,10 +133,15 @@ ROLE_ACTIONS: dict[str, tuple[tuple[str, ...], tuple[str, ...]]] = {
# Body phrases that prove an issue is an implementation container, not a # Body phrases that prove an issue is an implementation container, not a
# unit of direct author work (#844). Matched case-insensitively against the # unit of direct author work (#844 / #854). Matched case-insensitively against
# issue body. Title alone is never sufficient (ordinary issues may mention # the issue body. Title alone is never sufficient (ordinary issues may mention
# "epic" incidentally). # "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.
_CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = ( _CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
# Epic / child-only (#844, live #631)
"implementation is delivered via child issues only", "implementation is delivered via child issues only",
"implementation is delivered through child issues only", "implementation is delivered through child issues only",
"implementation is delivered via child issues", "implementation is delivered via child issues",
@@ -150,9 +155,26 @@ _CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
"coordination container", "coordination container",
"child-only container", "child-only container",
"implementation is delegated to child", "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 labels (structured evidence preferred over title). # 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.
_EPIC_LABELS: frozenset[str] = frozenset( _EPIC_LABELS: frozenset[str] = frozenset(
{ {
"type:epic", "type:epic",
@@ -161,6 +183,16 @@ _EPIC_LABELS: frozenset[str] = frozenset(
"scope:epic", "scope:epic",
"type:umbrella", "type:umbrella",
"umbrella", "umbrella",
"kind:umbrella",
"scope:umbrella",
"type:vision",
"vision",
"kind:vision",
"scope:vision",
"type:roadmap",
"roadmap",
"kind:roadmap",
"scope:roadmap",
} }
) )
@@ -222,16 +254,45 @@ 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( def classify_epic_or_child_only_container(
c: WorkCandidate, c: WorkCandidate,
) -> tuple[bool, str | None]: ) -> tuple[bool, str | None]:
"""Return whether *c* is an epic / child-only implementation container (#844). """Return whether *c* is a non-implementable coordination container (#844/#854).
Exclusion uses structured evidence first (labels, body scope language). Exclusion uses structured evidence first (labels, body scope language).
A bare title containing the word "epic" is **not** enough — ordinary A bare title containing the words "epic", "roadmap", "vision", or
implementable issues may mention epics incidentally. A title that is "umbrella" is **not** enough — ordinary implementable issues may mention
explicitly prefixed ``Epic:`` only counts when the body also proves those terms incidentally. Explicit title prefixes (``Epic:``, ``Roadmap:``,
child-only / no-direct-implementation scope (or an epic label is present). ``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.
PRs are never classified as containers here (they already have a head). PRs are never classified as containers here (they already have a head).
""" """
@@ -245,25 +306,28 @@ def classify_epic_or_child_only_container(
title_l = title.lower() title_l = title.lower()
body_hits = [m for m in _CHILD_ONLY_BODY_MARKERS if m in body_l] body_hits = [m for m in _CHILD_ONLY_BODY_MARKERS if m in body_l]
title_epic_prefix = title_l.startswith("epic:") or title_l.startswith("epic ") title_prefix = _title_container_prefix(title_l)
if epic_label: if epic_label:
detail = f"label={epic_label[0]}" detail = f"label={epic_label[0]}"
if body_hits: if body_hits:
detail = f"{detail}; body_marker={body_hits[0]!r}" detail = f"{detail}; body_marker={body_hits[0]!r}"
if title_prefix:
detail = f"{title_prefix}; {detail}"
return True, detail return True, detail
if body_hits: if body_hits:
# Body proves child-only / umbrella scope. Title "Epic:" is corroborating # Body proves child-only / vision / roadmap / umbrella scope. Title
# but not required — containers without the word still exclude. # prefixes are corroborating but not required — containers without the
# title word still exclude.
detail = f"body_marker={body_hits[0]!r}" detail = f"body_marker={body_hits[0]!r}"
if title_epic_prefix: if title_prefix:
detail = f"title_epic_prefix; {detail}" detail = f"{title_prefix}; {detail}"
return True, detail return True, detail
# Title-only "Epic:" without body scope evidence is insufficient (#844 AC: # Title-only coordination prefix without body scope evidence is
# eligibility does not rely solely on the word "Epic" in a title). # insufficient (#844/#854 AC: eligibility does not rely solely on a title
# Similarly, incidental "epic" mid-title without markers stays eligible. # word). Incidental mid-title mentions without markers stay eligible.
return False, None return False, None
+54 -6
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) |
@@ -75,7 +77,9 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| `/api/actions/{id}/preview` | Mutation ledger preview (GET, read-only) | | `/api/actions/{id}/preview` | Mutation ledger preview (GET, read-only) |
| `/leases` | Lease and collision visibility (#433) | | `/leases` | Lease and collision visibility (#433) |
| `/api/leases` | JSON lease/collision export | | `/api/leases` | JSON lease/collision export |
| `/sessions` | Phase 1 shell stub — session inventory (backed by #636) | | `/sessions` | Runtime and session view (#641) — health + inventory sessions/namespaces/worktrees |
| `/api/sessions` | JSON export for the runtime/session view |
| `/api/v1/sessions` | Versioned alias of `/api/sessions` |
| `/inventory` | Phase 1 shell stub — unified inventory (backed by #636) | | `/inventory` | Phase 1 shell stub — unified inventory (backed by #636) |
| `/timeline` | Phase 1 shell stub — workflow event timeline | | `/timeline` | Phase 1 shell stub — workflow event timeline |
| `/policy` | Phase 1 shell stub — capability/role policy placeholder | | `/policy` | Phase 1 shell stub — capability/role policy placeholder |
@@ -85,6 +89,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
@@ -253,11 +286,26 @@ The header carries two read-only status badges — an **environment** badge
a **mode: read-only** badge — plus a **Docs** link to this document. No a **mode: read-only** badge — plus a **Docs** link to this document. No
privileged action controls are present in the Phase 1 shell. privileged action controls are present in the Phase 1 shell.
Not-yet-implemented surfaces (`/sessions`, `/inventory`, `/timeline`, Not-yet-implemented surfaces (`/inventory`, `/timeline`, `/policy`,
`/policy`, `/insights`) resolve to graceful read-only stub pages instead of `/insights`) resolve to graceful read-only stub pages instead of 404s; their
404s; their backing views land in later child issues of #631 (the inventory backing views land in later child issues of #631 (the inventory surfaces are
surfaces are backed by #636). Mutating methods on stub routes still fail closed backed by #636). Mutating methods on stub routes still fail closed with
with `read-only-mvp`. `read-only-mvp`.
### Runtime and sessions (#641)
`/sessions` is a live Phase 1 read-only view that composes:
* runtime health from `#430` (profile, role, identity, master parity, stale warning)
* control-plane sessions / leases and filesystem locks / worktrees / namespaces from `#636`
* durable contamination markers when detectable (`#630` runtime recovery, `#671` stable-branch push)
It surfaces stale indicators (dead PID, expired lease) and never silences an
active contamination marker. Recovery links point only at sanctioned
reconnect/operator restart docs (`docs/mcp-namespace-eof-recovery.md`,
`docs/mcp-namespace-health.md`, `docs/mcp-restart-path-inventory.md`, this
document). The page does **not** restart, kill, or take over sessions; manual
`pkill` of MCP daemons is contamination, not recovery.
## System-health dashboard (#639) ## System-health dashboard (#639)
@@ -0,0 +1,378 @@
"""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()
+457
View File
@@ -0,0 +1,457 @@
"""Tests for the Runtime and session view (Phase 1, #641).
Covers clean and stale session rendering, contamination marker surfacing,
worktree binding display, sanctioned recovery links (no pkill), nav/live
status, and the JSON API export.
"""
from __future__ import annotations
import sys
import unittest
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from tests.webui_testclient import TestClient
from webui.app import create_app
from webui.inventory import (
AUTHORITY_CONTROL_PLANE_DB,
AUTHORITY_FILESYSTEM,
InventorySection,
InventorySnapshot,
STATUS_OK,
STATUS_UNAVAILABLE,
)
from webui.nav import NAV_GROUPS, STUB_PAGES, iter_nav_items
from webui.runtime_health import FileHash, RuntimeSnapshot
from webui.session_loader import (
ContaminationMarker,
SessionRow,
SessionViewSnapshot,
_build_session_rows,
load_session_view_snapshot,
snapshot_to_dict,
)
from webui.session_views import render_sessions_page
def _runtime(
*,
stale: str | None = None,
profile: str = "prgs-author",
role: str = "author",
) -> RuntimeSnapshot:
return RuntimeSnapshot(
project_id="gitea-tools",
repo_root="/tmp/repo",
remote="prgs",
host="gitea.prgs.cc",
profile_name=profile,
role_kind=role,
config_model="v2-contexts",
profile_mode="dynamic-profile",
profile_source="config file profile",
authenticated_username="jcwalker3",
identity_error=None,
repo_sha="a" * 40,
remote_master_sha="a" * 40,
commits_behind_master=0,
stale_runtime_warning=stale,
shell_health={"shell_use_allowed": True, "consecutive_spawn_failures": 0},
workflow_hashes=(
FileHash(label="SKILL.md", path="skills/llm-project-workflow/SKILL.md", sha256="abc"),
),
schema_hashes=(),
restart_guidance="docs/mcp-namespace-eof-recovery.md",
fetch_error=None,
)
def _inventory(
*,
sessions: tuple[dict, ...] = (),
leases: tuple[dict, ...] = (),
locks: tuple[dict, ...] = (),
worktrees: tuple[dict, ...] = (),
namespaces: tuple[dict, ...] = (),
) -> InventorySnapshot:
sections = (
InventorySection(
name="sessions",
authority=AUTHORITY_CONTROL_PLANE_DB,
status=STATUS_OK,
items=sessions,
),
InventorySection(
name="leases",
authority=AUTHORITY_CONTROL_PLANE_DB,
status=STATUS_OK,
items=leases,
),
InventorySection(
name="locks",
authority=AUTHORITY_FILESYSTEM,
status=STATUS_OK,
items=locks,
),
InventorySection(
name="worktrees",
authority=AUTHORITY_FILESYSTEM,
status=STATUS_OK,
items=worktrees,
),
InventorySection(
name="namespaces",
authority=AUTHORITY_FILESYSTEM,
status=STATUS_OK,
items=namespaces
or (
{
"profile_name": "prgs-author",
"role": "author",
"mcp_namespace": "gitea-author",
"capability_summary": {
"can_author": True,
"can_review": False,
"can_merge": False,
},
"active": True,
},
),
reason="only the profile serving this web process is observable",
),
)
index = {section.name: section for section in sections}
return InventorySnapshot(
generated_at="2026-07-25T00:00:00+00:00",
sections=sections,
collisions=(),
correlations=(),
scan_ms=1.0,
_section_index=index,
)
def _clean_session() -> dict:
return {
"session_id": "prgs-author-111-clean",
"role": "author",
"profile": "prgs-author",
"namespace": "gitea-author",
"pid": 1111,
"pid_alive": True,
"status": "active",
"started_at": "2026-07-25T00:00:00Z",
"last_heartbeat_at": "2026-07-25T01:00:00Z",
}
def _stale_session() -> dict:
return {
"session_id": "prgs-author-222-stale",
"role": "author",
"profile": "prgs-author",
"namespace": "gitea-author",
"pid": 2222,
"pid_alive": False,
"status": "active",
"started_at": "2026-07-24T00:00:00Z",
"last_heartbeat_at": "2026-07-24T01:00:00Z",
}
class TestBuildSessionRows(unittest.TestCase):
def test_clean_session_has_no_stale_or_contamination_flags(self):
inventory = _inventory(
sessions=(_clean_session(),),
leases=(
{
"lease_id": "lease-clean",
"session_id": "prgs-author-111-clean",
"status": "active",
"expired": False,
"work_kind": "issue",
"work_number": 641,
},
),
locks=(
{
"issue_number": 641,
"branch_name": "feat/issue-641-runtime-session-view",
"worktree_path": "~/Development/Gitea-Tools/branches/feat-issue-641",
"live": True,
},
),
)
rows = _build_session_rows(inventory, contamination=())
self.assertEqual(len(rows), 1)
row = rows[0]
self.assertEqual(row.session_id, "prgs-author-111-clean")
self.assertEqual(row.role, "author")
self.assertEqual(row.namespace, "gitea-author")
self.assertEqual(row.pid_alive, True)
self.assertEqual(row.lease_ids, ("lease-clean",))
self.assertEqual(row.work_refs, ("issue#641",))
self.assertTrue(row.worktree_paths)
self.assertEqual(row.stale_flags, ())
self.assertEqual(row.contamination_flags, ())
def test_stale_session_flags_dead_pid(self):
inventory = _inventory(sessions=(_stale_session(),))
rows = _build_session_rows(inventory, contamination=())
self.assertEqual(rows[0].stale_flags, ("pid-dead",))
def test_contamination_marker_binds_to_session(self):
inventory = _inventory(sessions=(_clean_session(),))
marker = ContaminationMarker(
kind="runtime_recovery_contamination",
on_disk=True,
has_payload=True,
summary="manual daemon kill",
reason_class="manual_daemon_kill",
session_id="prgs-author-111-clean",
role="author",
command_summary="pkill -f mcp_server.py",
cleared=False,
)
rows = _build_session_rows(inventory, contamination=(marker,))
self.assertIn("runtime_recovery_contamination", rows[0].contamination_flags)
def test_process_wide_contamination_surfaces_on_all_sessions(self):
inventory = _inventory(sessions=(_clean_session(), _stale_session()))
marker = ContaminationMarker(
kind="stable_branch_contamination",
on_disk=True,
has_payload=True,
summary="direct master push attempt",
reason_class="stable_branch_push",
session_id=None,
cleared=False,
)
rows = _build_session_rows(inventory, contamination=(marker,))
self.assertEqual(len(rows), 2)
for row in rows:
self.assertTrue(
any("stable_branch_contamination" in f for f in row.contamination_flags)
)
class TestRenderSessionsPage(unittest.TestCase):
def _snapshot(
self,
*,
sessions: tuple[dict, ...],
contamination: tuple[ContaminationMarker, ...] = (),
stale_runtime: str | None = None,
) -> SessionViewSnapshot:
inventory = _inventory(
sessions=sessions,
leases=(
{
"lease_id": "lease-1",
"session_id": sessions[0]["session_id"] if sessions else "",
"status": "active",
"expired": False,
"work_kind": "issue",
"work_number": 641,
},
)
if sessions
else (),
locks=(
{
"issue_number": 641,
"worktree_path": "branches/feat-issue-641",
},
)
if sessions
else (),
worktrees=(
{
"rel_path": "branches/feat-issue-641",
"branch": "feat/issue-641-runtime-session-view",
"classification": "active_issue_work",
"registered_worktree": True,
"dirty": False,
},
),
)
rows = _build_session_rows(inventory, contamination)
return SessionViewSnapshot(
runtime=_runtime(stale=stale_runtime),
inventory=inventory,
sessions=rows,
contamination_markers=contamination,
)
def test_clean_session_render(self):
html = render_sessions_page(self._snapshot(sessions=(_clean_session(),)))
self.assertIn("Runtime and sessions", html)
self.assertIn("prgs-author-111-clean", html)
self.assertIn("gitea-author", html)
self.assertIn("branches/feat-issue-641", html)
self.assertIn("Sanctioned recovery", html)
self.assertIn("docs/mcp-namespace-eof-recovery.md", html)
# Recovery section must name reconnect and forbid manual kill.
recovery_idx = html.lower().find("sanctioned recovery")
self.assertGreaterEqual(recovery_idx, 0)
recovery = html[recovery_idx:].lower()
self.assertIn("reconnect", recovery)
self.assertIn("contamination", recovery)
self.assertIn("not recovery", recovery)
self.assertNotIn("run pkill", recovery)
self.assertNotIn("killall", recovery)
def test_stale_session_render(self):
html = render_sessions_page(self._snapshot(sessions=(_stale_session(),)))
self.assertIn("prgs-author-222-stale", html)
self.assertIn("pid-dead", html)
self.assertIn("badge-stale", html)
def test_contamination_render_is_not_silent(self):
marker = ContaminationMarker(
kind="runtime_recovery_contamination",
on_disk=True,
has_payload=True,
summary="manual kill",
reason_class="manual_daemon_kill",
session_id="prgs-author-111-clean",
command_summary="pkill -f mcp_server.py",
cleared=False,
)
html = render_sessions_page(
self._snapshot(sessions=(_clean_session(),), contamination=(marker,))
)
self.assertIn("Contamination markers", html)
self.assertIn("runtime_recovery_contamination", html)
self.assertIn("ACTIVE", html)
self.assertIn("badge-blocked", html)
def test_stale_runtime_banner(self):
html = render_sessions_page(
self._snapshot(
sessions=(_clean_session(),),
stale_runtime="server behind master by 3 commits",
)
)
self.assertIn("Stale runtime", html)
self.assertIn("server behind master", html)
class TestSessionLoaderComposition(unittest.TestCase):
def test_load_with_injected_sources(self):
inventory = _inventory(sessions=(_clean_session(), _stale_session()))
snap = load_session_view_snapshot(
load_runtime=lambda: _runtime(),
load_inventory=lambda: inventory,
inspect_contamination=lambda **_k: {
"on_disk": False,
"has_payload": False,
"summary": "absent",
},
load_contamination_payload=lambda **_k: None,
)
self.assertEqual(len(snap.sessions), 2)
self.assertEqual(snap.stale_session_count, 1)
self.assertEqual(snap.contaminated_session_count, 0)
data = snapshot_to_dict(snap)
self.assertEqual(data["view"], "runtime-sessions")
self.assertEqual(data["issue"], 641)
self.assertTrue(data["read_only"])
self.assertEqual(data["session_counts"]["total"], 2)
self.assertEqual(data["session_counts"]["stale"], 1)
self.assertIn("recovery_docs", data)
class TestSessionsRoutes(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
inventory = _inventory(
sessions=(_clean_session(), _stale_session()),
leases=(
{
"lease_id": "lease-x",
"session_id": "prgs-author-111-clean",
"status": "active",
"expired": False,
"work_kind": "issue",
"work_number": 641,
},
),
locks=(
{
"issue_number": 641,
"worktree_path": "branches/feat-issue-641",
},
),
worktrees=(
{
"rel_path": "branches/feat-issue-641",
"branch": "feat/issue-641-runtime-session-view",
"classification": "active_issue_work",
"registered_worktree": True,
"dirty": False,
},
),
)
rows = _build_session_rows(inventory, contamination=())
self.snapshot = SessionViewSnapshot(
runtime=_runtime(stale="stale for test"),
inventory=inventory,
sessions=rows,
contamination_markers=(),
)
self._patch = mock.patch(
"webui.app.load_session_view_snapshot",
return_value=self.snapshot,
)
self._patch.start()
def tearDown(self):
self._patch.stop()
def test_sessions_page_live(self):
response = self.client.get("/sessions")
self.assertEqual(response.status_code, 200)
self.assertIn("Runtime and sessions", response.text)
self.assertIn("prgs-author-111-clean", response.text)
self.assertIn("prgs-author-222-stale", response.text)
self.assertIn("pid-dead", response.text)
self.assertIn("Sanctioned recovery", response.text)
self.assertNotIn("Phase 1 shell placeholder", response.text)
self.assertNotIn("child issue of #425", response.text.lower())
def test_api_sessions_json(self):
for path in ("/api/sessions", "/api/v1/sessions"):
response = self.client.get(path)
self.assertEqual(response.status_code, 200, path)
data = response.json()
self.assertEqual(data["view"], "runtime-sessions")
self.assertEqual(data["session_counts"]["total"], 2)
self.assertEqual(data["session_counts"]["stale"], 1)
self.assertTrue(data["read_only"])
self.assertEqual(data["mutations"], [])
def test_nav_marks_sessions_live(self):
sessions_items = [
item for item in iter_nav_items() if item.href == "/sessions"
]
self.assertEqual(len(sessions_items), 1)
self.assertEqual(sessions_items[0].status, "live")
self.assertNotIn("/sessions", STUB_PAGES)
# Home page should not mark Sessions as stub.
home = self.client.get("/")
self.assertEqual(home.status_code, 200)
self.assertIn('href="/sessions"', home.text)
# Stub marker only appears next to remaining stub destinations.
self.assertNotIn(
'href="/sessions">Sessions</a> <span class="muted">(stub)</span>',
home.text,
)
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()
+34
View File
@@ -42,10 +42,17 @@ 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
from webui.runtime_views import render_runtime_page from webui.runtime_views import render_runtime_page
from webui.session_loader import (
load_session_view_snapshot,
snapshot_to_dict as session_view_snapshot_to_dict,
)
from webui.session_views import render_sessions_page
from webui.inventory import ( from webui.inventory import (
SECTION_NAMES as _INVENTORY_SECTIONS, SECTION_NAMES as _INVENTORY_SECTIONS,
load_inventory_snapshot, load_inventory_snapshot,
@@ -80,10 +87,12 @@ 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)"),
("/runtime", "Runtime", "MCP health and stale-runtime detection (#430)"), ("/runtime", "Runtime", "MCP health and stale-runtime detection (#430)"),
("/sessions", "Sessions", "runtime and session view (#641)"),
("/audit", "Audit", "final-report paste and validator preview (#431)"), ("/audit", "Audit", "final-report paste and validator preview (#431)"),
("/worktrees", "Worktrees", "branch hygiene dashboard (#432)"), ("/worktrees", "Worktrees", "branch hygiene dashboard (#432)"),
("/leases", "Leases", "collision and lease visibility (#433)"), ("/leases", "Leases", "collision and lease visibility (#433)"),
@@ -200,6 +209,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:
@@ -313,6 +331,17 @@ async def api_runtime(_request: Request) -> JSONResponse:
return JSONResponse(runtime_snapshot_to_dict(load_runtime_snapshot())) return JSONResponse(runtime_snapshot_to_dict(load_runtime_snapshot()))
async def sessions(_request: Request) -> HTMLResponse:
"""Runtime and session view (#641) — read-only composition of health + inventory."""
snapshot = load_session_view_snapshot()
return HTMLResponse(render_sessions_page(snapshot))
async def api_sessions(_request: Request) -> JSONResponse:
"""JSON export for the runtime/session view (#641)."""
return JSONResponse(session_view_snapshot_to_dict(load_session_view_snapshot()))
async def _parse_audit_form(request: Request) -> tuple[str, str | None]: async def _parse_audit_form(request: Request) -> tuple[str, str | None]:
if request.method == "GET": if request.method == "GET":
return "", None return "", None
@@ -736,6 +765,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"]),
@@ -750,6 +781,9 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/api/prompts", api_prompts, methods=["GET"]), Route("/api/prompts", api_prompts, methods=["GET"]),
Route("/runtime", runtime, methods=["GET"]), Route("/runtime", runtime, methods=["GET"]),
Route("/api/runtime", api_runtime, methods=["GET"]), Route("/api/runtime", api_runtime, methods=["GET"]),
Route("/sessions", sessions, methods=["GET"]),
Route("/api/sessions", api_sessions, methods=["GET"]),
Route("/api/v1/sessions", api_sessions, methods=["GET"]),
Route("/api/v1/timeline", api_v1_timeline, methods=["GET"]), Route("/api/v1/timeline", api_v1_timeline, methods=["GET"]),
Route("/analytics", analytics, methods=["GET"]), Route("/analytics", analytics, methods=["GET"]),
Route("/api/analytics", api_v1_analytics, methods=["GET"]), Route("/api/analytics", api_v1_analytics, 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"),
+2 -6
View File
@@ -41,13 +41,14 @@ 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"),
)), )),
NavGroup("Runtime/Sessions", ( NavGroup("Runtime/Sessions", (
NavItem("/runtime", "Runtime health"), NavItem("/runtime", "Runtime health"),
NavItem("/sessions", "Sessions", "stub"), NavItem("/sessions", "Sessions"),
)), )),
NavGroup("Projects", ( NavGroup("Projects", (
NavItem("/projects", "Projects"), NavItem("/projects", "Projects"),
@@ -75,11 +76,6 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
# issues of epic #631. Each maps a path to (title, description). Routes are # issues of epic #631. Each maps a path to (title, description). Routes are
# registered so nav links resolve to a graceful, read-only stub page. # registered so nav links resolve to a graceful, read-only stub page.
STUB_PAGES: dict[str, tuple[str, str]] = { STUB_PAGES: dict[str, tuple[str, str]] = {
"/sessions": (
"Sessions",
"Active session, capability, and role inventory. Backed by the unified "
"inventory API (#636) once it lands.",
),
"/inventory": ( "/inventory": (
"Inventory", "Inventory",
"Unified sessions, leases, locks, namespaces, and worktree inventory. " "Unified sessions, leases, locks, namespaces, and worktree inventory. "
+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 ""),
},
) )
+3 -1
View File
@@ -88,5 +88,7 @@ def render_runtime_page(snapshot: RuntimeSnapshot) -> str:
"<p class='muted'>MVP is read-only — restart MCP servers from your IDE/operator " "<p class='muted'>MVP is read-only — restart MCP servers from your IDE/operator "
"workflow. Related issue: <code>#420</code>. Guidance: " "workflow. Related issue: <code>#420</code>. Guidance: "
f"<code>{html.escape(snapshot.restart_guidance)}</code></p>" f"<code>{html.escape(snapshot.restart_guidance)}</code></p>"
"<p class='muted'>This page does not expose tokens or perform MCP restarts.</p>" "<p class='muted'>This page does not expose tokens or perform MCP restarts. "
"Correlated sessions, worktree bindings, and contamination markers: "
"<a href='/sessions'>/sessions</a> (#641).</p>"
) )
+428
View File
@@ -0,0 +1,428 @@
"""Compose runtime health + inventory into a sessions/runtime view (#641).
Phase 1 is read-only. It correlates namespaces, sessions, capabilities,
worktree bindings, lease ownership, stale flags, and contamination markers
when they are detectable on disk (#630 / #671). It never restarts, kills, or
takes over a session.
Sources:
* :mod:`webui.runtime_health` — profile, role, stale runtime, shell health.
* :mod:`webui.inventory` — sessions, leases, locks, worktrees, namespaces
and collision signals from the control-plane DB + filesystem.
* :mod:`mcp_session_state` — durable contamination markers (inspect only).
Secrets are never read. Absolute paths are collapsed by inventory redaction.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from typing import Any, Callable
import mcp_session_state
from webui.inventory import (
InventorySnapshot,
load_inventory_snapshot,
snapshot_to_dict as inventory_snapshot_to_dict,
)
from webui.runtime_health import (
RuntimeSnapshot,
load_runtime_snapshot,
snapshot_to_dict as runtime_snapshot_to_dict,
)
# Sanctioned recovery pointers only — never pkill / killall (#630).
SANCTIONED_RECOVERY_DOCS: tuple[dict[str, str], ...] = (
{
"label": "MCP namespace EOF recovery (reconnect only)",
"path": "docs/mcp-namespace-eof-recovery.md",
"note": "IDE/client reconnect or operator-owned restart; never kill daemons.",
},
{
"label": "MCP namespace health",
"path": "docs/mcp-namespace-health.md",
"note": "client_namespace probe proves namespace health.",
},
{
"label": "Restart path inventory",
"path": "docs/mcp-restart-path-inventory.md",
"note": "Catalog of sanctioned reconnect/restart paths.",
},
{
"label": "Local web UI recovery",
"path": "docs/webui-local-dev.md",
"note": "Operator console start and documented recovery sequence.",
},
)
_CONTAMINATION_KINDS: tuple[str, ...] = (
mcp_session_state.KIND_RUNTIME_RECOVERY_CONTAMINATION,
mcp_session_state.KIND_STABLE_BRANCH_CONTAMINATION,
)
@dataclass(frozen=True)
class ContaminationMarker:
"""A detectable durable contamination marker (audit-safe summary)."""
kind: str
on_disk: bool
has_payload: bool
summary: str
reason_class: str | None = None
session_id: str | None = None
role: str | None = None
command_summary: str | None = None
cleared: bool = False
def to_dict(self) -> dict[str, Any]:
return {
"kind": self.kind,
"on_disk": self.on_disk,
"has_payload": self.has_payload,
"summary": self.summary,
"reason_class": self.reason_class,
"session_id": self.session_id,
"role": self.role,
"command_summary": self.command_summary,
"cleared": self.cleared,
"active": self.on_disk and self.has_payload and not self.cleared,
}
@dataclass(frozen=True)
class SessionRow:
"""One correlated session row for the sessions table."""
session_id: str
role: str | None
profile: str | None
namespace: str | None
pid: int | None
pid_alive: bool | None
status: str | None
started_at: str | None
last_heartbeat_at: str | None
lease_ids: tuple[str, ...] = ()
work_refs: tuple[str, ...] = ()
worktree_paths: tuple[str, ...] = ()
stale_flags: tuple[str, ...] = ()
contamination_flags: tuple[str, ...] = ()
def to_dict(self) -> dict[str, Any]:
return {
"session_id": self.session_id,
"role": self.role,
"profile": self.profile,
"namespace": self.namespace,
"pid": self.pid,
"pid_alive": self.pid_alive,
"status": self.status,
"started_at": self.started_at,
"last_heartbeat_at": self.last_heartbeat_at,
"lease_ids": list(self.lease_ids),
"work_refs": list(self.work_refs),
"worktree_paths": list(self.worktree_paths),
"stale_flags": list(self.stale_flags),
"contamination_flags": list(self.contamination_flags),
"is_stale": bool(self.stale_flags),
"is_contaminated": bool(self.contamination_flags),
}
@dataclass(frozen=True)
class SessionViewSnapshot:
"""Composed runtime + session inventory view (#641)."""
runtime: RuntimeSnapshot
inventory: InventorySnapshot
sessions: tuple[SessionRow, ...]
contamination_markers: tuple[ContaminationMarker, ...]
recovery_docs: tuple[dict[str, str], ...] = SANCTIONED_RECOVERY_DOCS
fetch_error: str | None = None
@property
def stale_session_count(self) -> int:
return sum(1 for row in self.sessions if row.stale_flags)
@property
def contaminated_session_count(self) -> int:
return sum(1 for row in self.sessions if row.contamination_flags)
@property
def active_contamination(self) -> tuple[ContaminationMarker, ...]:
return tuple(m for m in self.contamination_markers if m.to_dict()["active"])
def _inspect_contamination(
kind: str,
*,
remote: str | None,
inspect: Callable[..., dict[str, Any]] | None = None,
load: Callable[..., dict[str, Any] | None] | None = None,
) -> ContaminationMarker:
"""Inspect one contamination kind; never raises into the page render path."""
inspect_fn = inspect or mcp_session_state.inspect_state_envelope
load_fn = load or mcp_session_state.load_state
try:
envelope = inspect_fn(kind=kind, remote=remote)
except Exception as exc: # noqa: BLE001 — fail soft for the dashboard
return ContaminationMarker(
kind=kind,
on_disk=False,
has_payload=False,
summary=f"contamination inspect failed: {type(exc).__name__}",
)
reason_class = None
session_id = None
role = None
command_summary = None
cleared = False
summary = str(envelope.get("summary") or "")
if envelope.get("on_disk") and envelope.get("has_payload"):
try:
payload = load_fn(kind=kind, remote=remote) or {}
except Exception: # noqa: BLE001
payload = {}
if isinstance(payload, dict):
reason_class = payload.get("reason_class")
session_id = payload.get("session_id")
role = payload.get("role")
command_summary = payload.get("command_summary") or payload.get("detail")
cleared = bool(payload.get("cleared_by_reconciler"))
if not summary:
summary = (
f"{kind}: {reason_class or 'present'}"
+ (" (cleared)" if cleared else "")
)
return ContaminationMarker(
kind=kind,
on_disk=bool(envelope.get("on_disk")),
has_payload=bool(envelope.get("has_payload")),
summary=summary or f"{kind}: not present",
reason_class=str(reason_class) if reason_class else None,
session_id=str(session_id) if session_id else None,
role=str(role) if role else None,
command_summary=str(command_summary) if command_summary else None,
cleared=cleared,
)
def _load_contamination_markers(
*,
remote: str | None,
inspect: Callable[..., dict[str, Any]] | None = None,
load: Callable[..., dict[str, Any] | None] | None = None,
) -> tuple[ContaminationMarker, ...]:
return tuple(
_inspect_contamination(kind, remote=remote, inspect=inspect, load=load)
for kind in _CONTAMINATION_KINDS
)
def _build_session_rows(
inventory: InventorySnapshot,
contamination: tuple[ContaminationMarker, ...],
) -> tuple[SessionRow, ...]:
sessions_section = inventory.section("sessions")
leases_section = inventory.section("leases")
locks_section = inventory.section("locks")
leases_by_session: dict[str, list[dict[str, Any]]] = {}
for lease in (leases_section.items if leases_section else ()):
sid = str(lease.get("session_id") or "")
if sid:
leases_by_session.setdefault(sid, []).append(lease)
locks_by_session_hint: dict[str, list[dict[str, Any]]] = {}
for lock in (locks_section.items if locks_section else ()):
# Locks carry claimant profile/username, not control-plane session ids.
# Correlate later by matching lease work_number ↔ lock issue_number.
pass
active_markers = [m for m in contamination if m.to_dict()["active"]]
marker_session_ids = {
m.session_id for m in active_markers if m.session_id
}
rows: list[SessionRow] = []
for raw in sessions_section.items if sessions_section else ():
sid = str(raw.get("session_id") or "")
if not sid:
continue
session_leases = leases_by_session.get(sid, [])
lease_ids = tuple(
str(lease["lease_id"])
for lease in session_leases
if lease.get("lease_id")
)
work_refs: list[str] = []
work_numbers: list[int] = []
for lease in session_leases:
kind = lease.get("work_kind")
number = lease.get("work_number")
if kind and number is not None:
work_refs.append(f"{kind}#{number}")
try:
work_numbers.append(int(number))
except (TypeError, ValueError):
pass
worktree_paths: list[str] = []
for lock in (locks_section.items if locks_section else ()):
try:
issue_no = int(lock.get("issue_number"))
except (TypeError, ValueError):
continue
if issue_no in work_numbers and lock.get("worktree_path"):
worktree_paths.append(str(lock["worktree_path"]))
stale_flags: list[str] = []
if raw.get("pid_alive") is False:
stale_flags.append("pid-dead")
status = str(raw.get("status") or "").lower()
if status and status not in {"active", "alive", "running", "ok"}:
stale_flags.append(f"status:{status}")
for lease in session_leases:
if lease.get("expired") is True:
stale_flags.append("lease-expired")
if str(lease.get("status") or "").lower() == "active" and lease.get(
"expired"
) is True:
stale_flags.append("active-lease-past-expiry")
contamination_flags: list[str] = []
if sid in marker_session_ids:
for marker in active_markers:
if marker.session_id == sid:
contamination_flags.append(marker.kind)
# Process-wide contamination with no session binding still surfaces
# against every live session so it cannot be silent (#630).
for marker in active_markers:
if not marker.session_id and marker.kind not in contamination_flags:
contamination_flags.append(f"{marker.kind}:process-wide")
rows.append(
SessionRow(
session_id=sid,
role=raw.get("role"),
profile=raw.get("profile"),
namespace=raw.get("namespace"),
pid=raw.get("pid") if isinstance(raw.get("pid"), int) else None,
pid_alive=raw.get("pid_alive")
if isinstance(raw.get("pid_alive"), bool)
else None,
status=raw.get("status"),
started_at=raw.get("started_at"),
last_heartbeat_at=raw.get("last_heartbeat_at"),
lease_ids=lease_ids,
work_refs=tuple(work_refs),
worktree_paths=tuple(worktree_paths),
stale_flags=tuple(dict.fromkeys(stale_flags)),
contamination_flags=tuple(dict.fromkeys(contamination_flags)),
)
)
# Silence unused variable for the locks_by_session_hint placeholder path.
_ = locks_by_session_hint
return tuple(rows)
def load_session_view_snapshot(
*,
load_runtime: Callable[..., RuntimeSnapshot] | None = None,
load_inventory: Callable[..., InventorySnapshot] | None = None,
inspect_contamination: Callable[..., dict[str, Any]] | None = None,
load_contamination_payload: Callable[..., dict[str, Any] | None] | None = None,
) -> SessionViewSnapshot:
"""Build the composed sessions/runtime view. Fail-soft on partial sources."""
runtime_loader = load_runtime or load_runtime_snapshot
inventory_loader = load_inventory or load_inventory_snapshot
fetch_error: str | None = None
try:
runtime = runtime_loader()
except Exception as exc: # noqa: BLE001
fetch_error = f"runtime snapshot failed: {type(exc).__name__}: {exc}"
# Minimal placeholder so the page still renders inventory + recovery.
from webui.runtime_health import RuntimeSnapshot as _RS
runtime = _RS(
project_id="unknown",
repo_root="",
remote="",
host="",
profile_name="unknown",
role_kind="unknown",
config_model="unknown",
profile_mode="unknown",
profile_source="unknown",
authenticated_username=None,
identity_error=str(exc),
repo_sha=None,
remote_master_sha=None,
commits_behind_master=None,
stale_runtime_warning=None,
shell_health={},
workflow_hashes=(),
schema_hashes=(),
restart_guidance="docs/mcp-namespace-eof-recovery.md",
fetch_error=str(exc),
)
try:
inventory = inventory_loader()
except Exception as exc: # noqa: BLE001
msg = f"inventory snapshot failed: {type(exc).__name__}: {exc}"
fetch_error = f"{fetch_error}; {msg}" if fetch_error else msg
inventory = load_inventory_snapshot(
db_path="/nonexistent-for-fail-soft",
lock_dir="/nonexistent-for-fail-soft",
)
remote = getattr(runtime, "remote", None)
contamination = _load_contamination_markers(
remote=remote,
inspect=inspect_contamination,
load=load_contamination_payload,
)
sessions = _build_session_rows(inventory, contamination)
return SessionViewSnapshot(
runtime=runtime,
inventory=inventory,
sessions=sessions,
contamination_markers=contamination,
fetch_error=fetch_error,
)
def snapshot_to_dict(snapshot: SessionViewSnapshot) -> dict[str, Any]:
"""JSON export for ``/api/sessions`` (read-only)."""
return {
"api_version": "v1",
"view": "runtime-sessions",
"issue": 641,
"fetch_error": snapshot.fetch_error,
"runtime": runtime_snapshot_to_dict(snapshot.runtime),
"inventory": inventory_snapshot_to_dict(snapshot.inventory),
"sessions": [row.to_dict() for row in snapshot.sessions],
"session_counts": {
"total": len(snapshot.sessions),
"stale": snapshot.stale_session_count,
"contaminated": snapshot.contaminated_session_count,
},
"contamination_markers": [
marker.to_dict() for marker in snapshot.contamination_markers
],
"active_contamination": [
marker.to_dict() for marker in snapshot.active_contamination
],
"recovery_docs": [dict(doc) for doc in snapshot.recovery_docs],
"read_only": True,
"phase": 1,
"mutations": [],
}
+344
View File
@@ -0,0 +1,344 @@
"""HTML views for the Runtime and session view (Phase 1, #641).
Read-only composition of runtime health (#430) and inventory sessions /
namespaces / worktrees (#636). Surfaces stale and contamination indicators
when detectable. Recovery links point only at sanctioned reconnect/restart
docs — never at manual process kill (#630).
"""
from __future__ import annotations
from html import escape
from typing import Sequence
from webui.layout import render_page
from webui.session_loader import (
ContaminationMarker,
SessionRow,
SessionViewSnapshot,
)
def _badge(text: str, css: str) -> str:
return f'<span class="badge {css}">{escape(text)}</span>'
def _flags(flags: Sequence[str], *, css: str) -> str:
if not flags:
return '<span class="muted">—</span>'
return " ".join(_badge(flag, css) for flag in flags)
def _runtime_banner(snapshot: SessionViewSnapshot) -> str:
runtime = snapshot.runtime
stale = runtime.stale_runtime_warning
stale_html = ""
if stale:
stale_html = (
f'<div class="health-card health-stale" style="margin-top:0.75rem;">'
f"<strong>Stale runtime:</strong> {escape(stale)}</div>"
)
identity = runtime.authenticated_username or "unresolved"
if runtime.identity_error:
identity = f"unresolved ({runtime.identity_error})"
return f"""<div class="health-card">
<h3>Runtime context</h3>
<table class="detail">
<tr><th>Profile</th><td><code>{escape(runtime.profile_name)}</code></td></tr>
<tr><th>Role kind</th><td>{escape(runtime.role_kind)}</td></tr>
<tr><th>Identity</th><td>{escape(str(identity))}</td></tr>
<tr><th>Remote / host</th>
<td><code>{escape(runtime.remote)}</code> · <code>{escape(runtime.host)}</code></td>
</tr>
<tr><th>Local HEAD</th>
<td><code>{escape(runtime.repo_sha or "unknown")}</code></td>
</tr>
<tr><th>Remote master</th>
<td><code>{escape(runtime.remote_master_sha or "unknown")}</code></td>
</tr>
<tr><th>Commits behind</th>
<td>{escape(str(runtime.commits_behind_master if runtime.commits_behind_master is not None else "unknown"))}</td>
</tr>
</table>
<p class="muted">Full runtime detail: <a href="/runtime">/runtime</a> ·
Inventory API: <a href="/api/v1/inventory"><code>/api/v1/inventory</code></a></p>
{stale_html}
</div>"""
def _summary_bar(snapshot: SessionViewSnapshot) -> str:
total = len(snapshot.sessions)
stale = snapshot.stale_session_count
contaminated = snapshot.contaminated_session_count
active_markers = len(snapshot.active_contamination)
inv_status = snapshot.inventory.status
return f"""<div class="health-card" style="display:flex; flex-wrap:wrap; gap:1rem; align-items:center;">
<div><strong>Sessions:</strong> <span class="badge badge-health-ok">{total}</span></div>
<div><strong>Stale:</strong> <span class="badge badge-stale">{stale}</span></div>
<div><strong>Contaminated:</strong> <span class="badge badge-blocked">{contaminated}</span></div>
<div><strong>Active markers:</strong> <span class="badge badge-health-unproven">{active_markers}</span></div>
<div><strong>Inventory:</strong> <span class="badge badge-health-skipped">{escape(inv_status)}</span></div>
</div>"""
def _render_session_row(row: SessionRow) -> str:
pid = "" if row.pid is None else str(row.pid)
pid_alive = "" if row.pid_alive is None else ("alive" if row.pid_alive else "dead")
pid_css = (
"badge-health-ok"
if row.pid_alive is True
else ("badge-blocked" if row.pid_alive is False else "badge-health-skipped")
)
leases = (
", ".join(f"<code>{escape(lid)}</code>" for lid in row.lease_ids)
if row.lease_ids
else '<span class="muted">none</span>'
)
work = (
", ".join(escape(ref) for ref in row.work_refs)
if row.work_refs
else '<span class="muted">—</span>'
)
worktrees = (
"<br>".join(f"<code>{escape(path)}</code>" for path in row.worktree_paths)
if row.worktree_paths
else '<span class="muted">unbound</span>'
)
return f"""<tr>
<td><code>{escape(row.session_id)}</code></td>
<td>
<div><code>{escape(str(row.role or ""))}</code> / <code>{escape(str(row.profile or ""))}</code></div>
<div class="muted" style="font-size:0.82rem;">ns: <code>{escape(str(row.namespace or ""))}</code></div>
</td>
<td>
<code>{escape(pid)}</code>
{_badge(pid_alive, pid_css)}
</td>
<td>{escape(str(row.status or ""))}<div class="muted" style="font-size:0.82rem;">{escape(str(row.last_heartbeat_at or ""))}</div></td>
<td>{leases}<div class="muted" style="font-size:0.82rem; margin-top:0.2rem;">{work}</div></td>
<td>{worktrees}</td>
<td>{_flags(row.stale_flags, css="badge-stale")}</td>
<td>{_flags(row.contamination_flags, css="badge-blocked")}</td>
</tr>"""
def _sessions_table(rows: Sequence[SessionRow]) -> str:
if not rows:
return (
'<p class="muted">No control-plane sessions recorded. Inventory may '
"be unavailable, or no MCP workers have registered yet.</p>"
)
body = "".join(_render_session_row(row) for row in rows)
return f"""<table class="registry">
<thead>
<tr>
<th>Session</th>
<th>Role / profile / namespace</th>
<th>PID</th>
<th>Status</th>
<th>Leases / work</th>
<th>Worktree binding</th>
<th>Stale</th>
<th>Contamination</th>
</tr>
</thead>
<tbody>
{body}
</tbody>
</table>"""
def _namespaces_section(snapshot: SessionViewSnapshot) -> str:
section = snapshot.inventory.section("namespaces")
if section is None:
return (
'<div class="prompt-card"><h3>Namespaces</h3>'
'<p class="muted">Namespaces section not loaded.</p></div>'
)
if not section.ok:
return f"""<div class="prompt-card">
<h3>Namespaces {_badge(section.status, "badge-health-degraded")}</h3>
<p class="muted">{escape(section.reason or "unavailable")}</p>
</div>"""
rows = []
for item in section.items:
caps = item.get("capability_summary") or {}
cap_bits = ", ".join(
name for name, ok in sorted(caps.items()) if ok
) or "none"
rows.append(
"<tr>"
f"<td><code>{escape(str(item.get('mcp_namespace') or ''))}</code></td>"
f"<td><code>{escape(str(item.get('profile_name') or ''))}</code></td>"
f"<td>{escape(str(item.get('role') or ''))}</td>"
f"<td>{escape(cap_bits)}</td>"
f"<td>{'yes' if item.get('active') else 'no'}</td>"
"</tr>"
)
reason = (
f'<p class="muted">{escape(section.reason)}</p>'
if section.reason
else ""
)
return f"""<div class="prompt-card">
<h3>Namespaces / capabilities</h3>
{reason}
<table class="registry">
<thead>
<tr>
<th>Namespace</th>
<th>Profile</th>
<th>Role</th>
<th>Capabilities</th>
<th>Active in process</th>
</tr>
</thead>
<tbody>
{"".join(rows) if rows else '<tr><td colspan="5" class="muted">No namespace rows.</td></tr>'}
</tbody>
</table>
</div>"""
def _worktrees_section(snapshot: SessionViewSnapshot) -> str:
section = snapshot.inventory.section("worktrees")
if section is None:
return ""
if not section.ok and not section.items:
return f"""<div class="prompt-card">
<h3>Worktrees {_badge(section.status, "badge-health-degraded")}</h3>
<p class="muted">{escape(section.reason or "unavailable")}</p>
</div>"""
rows = []
for item in section.items[:50]:
rows.append(
"<tr>"
f"<td><code>{escape(str(item.get('rel_path') or item.get('path') or ''))}</code></td>"
f"<td><code>{escape(str(item.get('branch') or ''))}</code></td>"
f"<td>{escape(str(item.get('classification') or ''))}</td>"
f"<td>{'yes' if item.get('registered_worktree') else 'no'}</td>"
f"<td>{'dirty' if item.get('dirty') else 'clean'}</td>"
"</tr>"
)
more = ""
if len(section.items) > 50:
more = f'<p class="muted">Showing 50 of {len(section.items)}. Full list: <a href="/worktrees">/worktrees</a>.</p>'
return f"""<div class="prompt-card">
<h3>Worktree bindings</h3>
<p class="muted">Registered issue worktrees under <code>branches/</code>. Hygiene detail: <a href="/worktrees">/worktrees</a>.</p>
<table class="registry">
<thead>
<tr>
<th>Path</th>
<th>Branch</th>
<th>Classification</th>
<th>Registered</th>
<th>State</th>
</tr>
</thead>
<tbody>
{"".join(rows) if rows else '<tr><td colspan="5" class="muted">No worktrees recorded.</td></tr>'}
</tbody>
</table>
{more}
</div>"""
def _contamination_section(markers: Sequence[ContaminationMarker]) -> str:
if not markers:
return (
'<div class="prompt-card"><h3>Contamination markers</h3>'
'<p class="muted">No contamination kinds inspected.</p></div>'
)
rows = []
for marker in markers:
active = marker.to_dict()["active"]
status = "ACTIVE" if active else ("cleared" if marker.cleared else "absent")
css = "badge-blocked" if active else "badge-health-ok"
rows.append(
"<tr>"
f"<td><code>{escape(marker.kind)}</code></td>"
f"<td>{_badge(status, css)}</td>"
f"<td>{escape(marker.reason_class or '')}</td>"
f"<td><code>{escape(marker.session_id or '')}</code></td>"
f"<td>{escape(marker.command_summary or marker.summary)}</td>"
"</tr>"
)
return f"""<div class="prompt-card">
<h3>Contamination markers (#630 / #671)</h3>
<p class="muted">Durable markers only — never silent when present. Clearance is reconciler-only.</p>
<table class="registry">
<thead>
<tr>
<th>Kind</th>
<th>State</th>
<th>Reason class</th>
<th>Session</th>
<th>Summary</th>
</tr>
</thead>
<tbody>
{"".join(rows)}
</tbody>
</table>
</div>"""
def _recovery_section(snapshot: SessionViewSnapshot) -> str:
items = []
for doc in snapshot.recovery_docs:
items.append(
"<li>"
f"<code>{escape(doc['path'])}</code> — "
f"<strong>{escape(doc['label'])}</strong>: {escape(doc['note'])}"
"</li>"
)
return f"""<div class="prompt-card">
<h3>Sanctioned recovery (read-only)</h3>
<p class="muted">This view does <strong>not</strong> restart, kill, or take over sessions.
Manual <code>pkill</code> / <code>kill</code> of MCP daemons is contamination (#630), not recovery.</p>
<ul class="reasons">
{"".join(items)}
<li>Prefer IDE/client reconnect (<code>/mcp reconnect</code>) or an operator-owned restart recorded in the restart inventory.</li>
</ul>
</div>"""
def render_sessions_page(snapshot: SessionViewSnapshot) -> str:
"""Render the full HTML body for the runtime/session view."""
error_block = ""
if snapshot.fetch_error:
error_block = (
f'<div class="health-card health-stale"><strong>Partial load:</strong> '
f"{escape(snapshot.fetch_error)}</div>"
)
if snapshot.runtime.fetch_error:
error_block += (
f'<div class="health-card health-stale"><strong>Runtime note:</strong> '
f"{escape(snapshot.runtime.fetch_error)}</div>"
)
body = f"""
{error_block}
{_runtime_banner(snapshot)}
{_summary_bar(snapshot)}
<div class="prompt-card">
<h3>Sessions</h3>
<p class="muted">Control-plane sessions correlated with leases and worktree bindings.
Stale and contamination flags are fail-soft: absence of a marker is not proof of cleanliness when inventory is degraded.</p>
{_sessions_table(snapshot.sessions)}
</div>
{_namespaces_section(snapshot)}
{_worktrees_section(snapshot)}
{_contamination_section(snapshot.contamination_markers)}
{_recovery_section(snapshot)}
"""
return render_page(title="Sessions", body_html=f"""<h2>Runtime and sessions</h2>
<p class="meta">Phase 1 read-only view (#641). Combines runtime health (#430) with
unified inventory sessions/namespaces/worktrees (#636). No restart or session-takeover controls.</p>
{body}""")
+4 -2
View File
@@ -278,8 +278,10 @@ def _recovery_card() -> str:
"controls arrive in Phase 2 (#642); until then recovery runs through " "controls arrive in Phase 2 (#642); until then recovery runs through "
"the sanctioned client reconnect / operator restart path.</p>" "the sanctioned client reconnect / operator restart path.</p>"
"<ul class='reasons'>" "<ul class='reasons'>"
"<li><a href='/runtime'>Runtime and session view</a> — active profile, " "<li><a href='/runtime'>Runtime health</a> — active profile, workflow "
"workflow hashes, and shell health.</li>" "hashes, and shell health.</li>"
"<li><a href='/sessions'>Runtime and sessions</a> — namespaces, session "
"rows, worktree bindings, and contamination markers (#641).</li>"
"<li>Reconnect the MCP client from the IDE, then re-run the blocked " "<li>Reconnect the MCP client from the IDE, then re-run the blocked "
"cycle. Never kill the daemon process manually: unmanaged kills are " "cycle. Never kill the daemon process manually: unmanaged kills are "
"recorded as runtime contamination (#630).</li>" "recorded as runtime contamination (#630).</li>"
+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)