Compare commits
14
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
930dc24632 | ||
|
|
7af40fb5ff | ||
|
|
3b68d15593 | ||
|
|
41622c5985 | ||
|
|
a4c73766f4 | ||
|
|
b2e28428a4 | ||
|
|
95e4aae287 | ||
|
|
301c78de20 | ||
|
|
dac40ab9b3 | ||
|
|
ccde9e8f11 | ||
|
|
1948d3dc21 | ||
|
|
714190e02a | ||
|
|
069a9af7e6 | ||
|
|
0a78da39e5 |
+81
-17
@@ -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
|
||||
# unit of direct author work (#844). Matched case-insensitively against the
|
||||
# issue body. Title alone is never sufficient (ordinary issues may mention
|
||||
# "epic" incidentally).
|
||||
# unit of direct author work (#844 / #854). Matched case-insensitively against
|
||||
# the issue body. Title alone is never sufficient (ordinary issues may mention
|
||||
# "epic", "roadmap", "vision", or "umbrella" incidentally).
|
||||
#
|
||||
# #854 extends the #844 marker set so product-vision (#652), phased-roadmap
|
||||
# (#653), and umbrella (#655) coordination records — which do not use the word
|
||||
# "epic" — are classified with the same semantic exclusion as epic containers.
|
||||
_CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
|
||||
# Epic / child-only (#844, live #631)
|
||||
"implementation is delivered via child issues only",
|
||||
"implementation is delivered through child issues only",
|
||||
"implementation is delivered via child issues",
|
||||
@@ -150,9 +155,26 @@ _CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
|
||||
"coordination container",
|
||||
"child-only container",
|
||||
"implementation is delegated to child",
|
||||
# Vision / roadmap / umbrella coordination (#854, live #652/#653/#655).
|
||||
# Prefer authoritative non-implementation / child-only scope language over
|
||||
# bare words like "roadmap" so ordinary implementable issues that mention
|
||||
# a parent vision or roadmap stay eligible.
|
||||
"do not implement features on this issue",
|
||||
"implementing features on this roadmap issue",
|
||||
"implementation is via linked children only",
|
||||
"no product feature claimed complete on this issue alone",
|
||||
"phased delivery roadmap and epic sequencing",
|
||||
"this issue is the enduring source of truth",
|
||||
"enduring source of truth for the",
|
||||
"canonical product vision — enduring source of truth",
|
||||
"canonical product vision - enduring source of truth",
|
||||
"state: vision-active",
|
||||
"state: roadmap-active",
|
||||
)
|
||||
|
||||
# Explicit epic / umbrella 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(
|
||||
{
|
||||
"type:epic",
|
||||
@@ -161,6 +183,16 @@ _EPIC_LABELS: frozenset[str] = frozenset(
|
||||
"scope:epic",
|
||||
"type:umbrella",
|
||||
"umbrella",
|
||||
"kind:umbrella",
|
||||
"scope:umbrella",
|
||||
"type:vision",
|
||||
"vision",
|
||||
"kind:vision",
|
||||
"scope:vision",
|
||||
"type:roadmap",
|
||||
"roadmap",
|
||||
"kind:roadmap",
|
||||
"scope:roadmap",
|
||||
}
|
||||
)
|
||||
|
||||
@@ -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(
|
||||
c: WorkCandidate,
|
||||
) -> 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).
|
||||
A bare title containing the word "epic" is **not** enough — ordinary
|
||||
implementable issues may mention epics incidentally. A title that is
|
||||
explicitly prefixed ``Epic:`` only counts when the body also proves
|
||||
child-only / no-direct-implementation scope (or an epic label is present).
|
||||
A bare title containing the words "epic", "roadmap", "vision", or
|
||||
"umbrella" is **not** enough — ordinary implementable issues may mention
|
||||
those terms incidentally. Explicit title prefixes (``Epic:``, ``Roadmap:``,
|
||||
``Product vision:``, ``Umbrella:``) only count when the body also proves
|
||||
child-only / no-direct-implementation scope (or a container label is
|
||||
present).
|
||||
|
||||
Covers epic, product-vision, phased-roadmap, umbrella, and child-only
|
||||
records so the allocator never assigns coordination containers as direct
|
||||
author work.
|
||||
|
||||
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()
|
||||
|
||||
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:
|
||||
detail = f"label={epic_label[0]}"
|
||||
if body_hits:
|
||||
detail = f"{detail}; body_marker={body_hits[0]!r}"
|
||||
if title_prefix:
|
||||
detail = f"{title_prefix}; {detail}"
|
||||
return True, detail
|
||||
|
||||
if body_hits:
|
||||
# Body proves child-only / umbrella scope. Title "Epic:" is corroborating
|
||||
# but not required — containers without the word still exclude.
|
||||
# Body proves child-only / vision / roadmap / umbrella scope. Title
|
||||
# prefixes are corroborating but not required — containers without the
|
||||
# title word still exclude.
|
||||
detail = f"body_marker={body_hits[0]!r}"
|
||||
if title_epic_prefix:
|
||||
detail = f"title_epic_prefix; {detail}"
|
||||
if title_prefix:
|
||||
detail = f"{title_prefix}; {detail}"
|
||||
return True, detail
|
||||
|
||||
# Title-only "Epic:" without body scope evidence is insufficient (#844 AC:
|
||||
# eligibility does not rely solely on the word "Epic" in a title).
|
||||
# Similarly, incidental "epic" mid-title without markers stays eligible.
|
||||
# Title-only coordination prefix without body scope evidence is
|
||||
# insufficient (#844/#854 AC: eligibility does not rely solely on a title
|
||||
# word). Incidental mid-title mentions without markers stay eligible.
|
||||
return False, None
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
# MCP restart classes and blast-radius permissions (#663)
|
||||
|
||||
This is the machine-enforced class matrix used by
|
||||
`restart_coordinator.RESTART_CLASS_POLICIES`. It implements the narrower-first
|
||||
recovery ladder from #655 and the authorization policy from #656, using the
|
||||
path inventory from #657 and the impact coordinator from #658. Product and
|
||||
delivery lineage: vision #652 and roadmap #653.
|
||||
|
||||
Unknown class names are denied. The coordinator requires both the class
|
||||
permission and an eligible request role. Approval gates are additional: a
|
||||
caller cannot turn a request permission into execution authority.
|
||||
|
||||
| Restart class | Required permission | Expected blast radius | Drain requirement | Approval requirement | Audit requirement | Recovery behavior |
|
||||
|---|---|---|---|---|---|---|
|
||||
| `client_reconnect` | `mcp.reconnect.client` | none | none | self service | class, actor, client namespace, reason, outcome | Reconnect only the caller's client transport. No daemon or peer work changes. |
|
||||
| `session_reconnect` | `mcp.reconnect.session` | low | requesting-session safe point | self service | class, actor, session, reason, outcome | Rebind identity, capability, and workspace state for one session. |
|
||||
| `worker_restart` | `mcp.restart.worker.request` | low | target worker | controller approval + automated gates | class, actor, worker, approval, scoped drain, outcome | Restart one worker after its own leases and mutations drain. |
|
||||
| `role_runtime_restart` | `mcp.restart.role_runtime.request` | medium | target role runtime | controller approval + automated gates | class, actor, role namespace, approval, scoped drain, outcome | Restart and re-probe one role runtime; unrelated roles remain available. |
|
||||
| `connector_restart` | `mcp.restart.connector.request` | medium | target connector | controller approval + automated gates | class, actor, connector, approval, scoped drain, outcome | Restart one connector while unrelated runtimes remain available. |
|
||||
| `configuration_reload` | `mcp.reload.configuration.request` | low | mutation quiesce | controller approval + automated gates | class, actor, configuration revision, approval, outcome | Gracefully reload configuration without replacing the daemon. |
|
||||
| `rolling_mcp_restart` | `mcp.restart.rolling.request` | medium | one instance at a time | controller approval + automated gates | class, actor, instance order, approval, per-instance drains, outcome | Drain, restart, verify, and restore each instance before advancing. |
|
||||
| `full_mcp_restart` | `mcp.restart.full.request` | high | all sessions and mutations | controller approval + automated gates | class, actor, full impact, approval, full drain proof, outcome | Replace the complete MCP runtime only after a verified full drain. |
|
||||
| `host_restart` | `mcp.restart.host.request` | high | all host work | controller approval + infrastructure operator | class, actor, host/change or incident id, approval, full drain proof, outcome | Hand off to infrastructure ownership and reconcile every runtime afterward. |
|
||||
|
||||
## Drain boundary
|
||||
|
||||
Only `full_mcp_restart` and `host_restart` set `full_drain_required=true`.
|
||||
Reconnects and configuration reloads do not disrupt peer sessions. Worker,
|
||||
role-runtime, and connector restarts evaluate only their explicitly named
|
||||
target. Rolling restart drains one instance at a time. Missing required target
|
||||
scope denies the request rather than silently widening it to a full restart.
|
||||
|
||||
## Permission and approval boundary
|
||||
|
||||
Author, reviewer, merger, and reconciler roles may self-request reconnects and
|
||||
request scoped worker/role/connector/reload recovery. They cannot request
|
||||
rolling, full, or host restart classes. Controller/operator/admin roles may
|
||||
request the broader classes, while execution remains operator/admin-owned.
|
||||
Controller approval is independently required for every class above a session
|
||||
reconnect. Host restart additionally requires infrastructure-operator proof.
|
||||
|
||||
The MCP request tool derives class permissions from its authenticated runtime
|
||||
role. It does not accept caller-supplied permissions. Controller and operator
|
||||
authorization are read from the already-running daemon environment, never
|
||||
from a request argument.
|
||||
|
||||
## Audit and failure behavior
|
||||
|
||||
Every impact audit and every console restart/reload audit includes a
|
||||
`restart_class` field. The impact audit also includes the exact
|
||||
`required_permission`. Unknown classes, missing permissions, ineligible roles,
|
||||
missing approval, missing scoped targets, and incomplete inventory all deny
|
||||
fail closed. Manual process kills remain forbidden and contaminating (#630).
|
||||
@@ -12,6 +12,12 @@ restart/reload/kill paths). The **mutative apply** path — actually performing
|
||||
restart — is a later child gated by a drain proof and is explicitly out of
|
||||
scope here.
|
||||
|
||||
The coordinator now routes every request through the restart-class policy
|
||||
matrix defined for #663. See
|
||||
[`mcp-restart-classes.md`](./mcp-restart-classes.md) for permissions, expected
|
||||
blast radius, scoped drain and approval requirements, audit fields, and
|
||||
recovery behavior for all nine classes.
|
||||
|
||||
## Components
|
||||
|
||||
| Piece | Where | Responsibility |
|
||||
@@ -77,7 +83,10 @@ authorization is present.
|
||||
```text
|
||||
gitea_request_mcp_restart(remote, host, org, repo,
|
||||
dry_run=True, request_override=False,
|
||||
session_id=None, limit=200)
|
||||
session_id=None, limit=200,
|
||||
restart_class="full_mcp_restart",
|
||||
target_session_id=None, target_role=None,
|
||||
target_connector=None)
|
||||
```
|
||||
|
||||
Read-only, dry-run, and it **never restarts anything**. `apply_supported` is
|
||||
@@ -87,7 +96,8 @@ apply is gated by a drain proof (a separate child).
|
||||
## Audit
|
||||
|
||||
Every evaluation carries an `audit_record` (event, coordinator version, verdict,
|
||||
allow decision, blast radius, counts, timestamp) so restart decisions are
|
||||
restart class, required permission, allow decision, blast radius, counts,
|
||||
timestamp) so restart decisions are
|
||||
auditable. No secrets flow through the coordinator — session ids, pids, and
|
||||
profiles are operational metadata only.
|
||||
|
||||
|
||||
@@ -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) |
|
||||
| `/queue` | Live PR and issue queue dashboard (#429) |
|
||||
| `/api/queue` | JSON queue export with pagination metadata |
|
||||
| `/traffic` | Workflow traffic-control view — runnable, leased, blocked, needs-controller, terminal-complete (#640) |
|
||||
| `/api/traffic` | JSON traffic-control export with state classifications and next safe role actions |
|
||||
| `/projects` | Project registry list with status and onboarding progress (#427, #635) |
|
||||
| `/projects/{id}` | Project detail + onboarding checklist |
|
||||
| `/api/v1/projects` | Versioned JSON registry export (#635) |
|
||||
@@ -85,6 +87,35 @@ Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with
|
||||
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
|
||||
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)
|
||||
|
||||
`GET /api/v1/system/health` is the structured, read-only health surface for
|
||||
|
||||
+29
-1
@@ -22342,10 +22342,15 @@ def gitea_request_mcp_restart(
|
||||
request_override: bool = False,
|
||||
session_id: str | None = None,
|
||||
limit: int = 200,
|
||||
restart_class: str = "full_mcp_restart",
|
||||
target_session_id: str | None = None,
|
||||
target_role: str | None = None,
|
||||
target_connector: str | None = None,
|
||||
) -> dict:
|
||||
"""Evaluate a proposed MCP restart and return an impact preview (#658).
|
||||
|
||||
Central restart coordinator: gathers live control-plane state (sessions,
|
||||
Central restart coordinator: resolves the requested restart class, gathers
|
||||
live control-plane state (sessions,
|
||||
leases/locks, in-flight issue/PR work, mutations, worktrees) and returns a
|
||||
blast-radius impact report with a ``safe`` / ``unsafe`` / ``override``
|
||||
verdict, so the console (#642/#652) and operators can see what a restart
|
||||
@@ -22439,6 +22444,9 @@ def gitea_request_mcp_restart(
|
||||
|
||||
profile = get_profile()
|
||||
profile_name = (profile.get("profile_name") or "").strip() or "session"
|
||||
requester_role = (
|
||||
profile.get("role_kind") or profile.get("role") or ""
|
||||
).strip().lower()
|
||||
sid = (session_id or "").strip() or f"{profile_name}-{os.getpid()}"
|
||||
|
||||
# Override authority is read from the environment only — a worker session
|
||||
@@ -22448,6 +22456,15 @@ def gitea_request_mcp_restart(
|
||||
(os.environ.get("GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION") or "").strip()
|
||||
)
|
||||
operator_override = bool(request_override and operator_authorized)
|
||||
controller_approved = bool(
|
||||
(
|
||||
os.environ.get("GITEA_CONTROLLER_RESTART_APPROVAL_AUTHORIZATION")
|
||||
or ""
|
||||
).strip()
|
||||
)
|
||||
requester_permissions = restart_coordinator.permissions_for_role(
|
||||
requester_role
|
||||
)
|
||||
|
||||
inventory = {
|
||||
"sessions": sessions,
|
||||
@@ -22462,6 +22479,14 @@ def gitea_request_mcp_restart(
|
||||
operator_override=operator_override,
|
||||
requesting_session_id=sid,
|
||||
dry_run=True, # coordinator is always analysis-only (#658)
|
||||
restart_class=restart_class,
|
||||
requester_role=requester_role,
|
||||
requester_permissions=requester_permissions,
|
||||
controller_approved=controller_approved,
|
||||
operator_authorized=operator_authorized,
|
||||
target_session_id=target_session_id,
|
||||
target_role=target_role,
|
||||
target_connector=target_connector,
|
||||
)
|
||||
|
||||
payload = report.as_dict()
|
||||
@@ -22473,6 +22498,9 @@ def gitea_request_mcp_restart(
|
||||
payload["requesting_session_id"] = sid
|
||||
payload["operator_override_requested"] = bool(request_override)
|
||||
payload["operator_override_authorized"] = operator_authorized
|
||||
payload["controller_approval_authorized"] = controller_approved
|
||||
payload["requester_role"] = requester_role
|
||||
payload["requester_permissions"] = list(requester_permissions)
|
||||
payload["apply_supported"] = False
|
||||
if not dry_run:
|
||||
payload["reasons"] = list(payload.get("reasons") or []) + [
|
||||
|
||||
+370
-14
@@ -28,11 +28,12 @@ from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from enum import Enum
|
||||
from typing import Any, Mapping, Sequence
|
||||
|
||||
import lease_lifecycle
|
||||
|
||||
COORDINATOR_VERSION = "1.0.0-issue-658"
|
||||
COORDINATOR_VERSION = "1.1.0-issue-663"
|
||||
|
||||
# Restart verdicts. Exactly the three the acceptance criteria name.
|
||||
VERDICT_SAFE = "safe"
|
||||
@@ -54,6 +55,194 @@ LEASE_FRESHNESS_LIVE = "active"
|
||||
DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS = 900
|
||||
|
||||
|
||||
class RestartClass(str, Enum):
|
||||
"""The only restart/recovery classes accepted by the coordinator."""
|
||||
|
||||
CLIENT_RECONNECT = "client_reconnect"
|
||||
SESSION_RECONNECT = "session_reconnect"
|
||||
WORKER_RESTART = "worker_restart"
|
||||
ROLE_RUNTIME_RESTART = "role_runtime_restart"
|
||||
CONNECTOR_RESTART = "connector_restart"
|
||||
CONFIGURATION_RELOAD = "configuration_reload"
|
||||
ROLLING_MCP_RESTART = "rolling_mcp_restart"
|
||||
FULL_MCP_RESTART = "full_mcp_restart"
|
||||
HOST_RESTART = "host_restart"
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class RestartClassPolicy:
|
||||
"""Least-privilege policy for one :class:`RestartClass`."""
|
||||
|
||||
restart_class: RestartClass
|
||||
required_permission: str
|
||||
expected_blast_radius: str
|
||||
drain_requirement: str
|
||||
full_drain_required: bool
|
||||
approval_requirement: str
|
||||
audit_requirement: str
|
||||
recovery_behavior: str
|
||||
request_roles: tuple[str, ...]
|
||||
execution_roles: tuple[str, ...]
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
"restart_class": self.restart_class.value,
|
||||
"required_permission": self.required_permission,
|
||||
"expected_blast_radius": self.expected_blast_radius,
|
||||
"drain_requirement": self.drain_requirement,
|
||||
"full_drain_required": self.full_drain_required,
|
||||
"approval_requirement": self.approval_requirement,
|
||||
"audit_requirement": self.audit_requirement,
|
||||
"recovery_behavior": self.recovery_behavior,
|
||||
"request_roles": list(self.request_roles),
|
||||
"execution_roles": list(self.execution_roles),
|
||||
}
|
||||
|
||||
|
||||
WORKER_ROLES = ("author", "reviewer", "merger", "reconciler")
|
||||
CONTROL_ROLES = ("controller", "operator", "admin")
|
||||
ALL_REQUEST_ROLES = WORKER_ROLES + CONTROL_ROLES
|
||||
|
||||
RESTART_CLASS_POLICIES: dict[RestartClass, RestartClassPolicy] = {
|
||||
RestartClass.CLIENT_RECONNECT: RestartClassPolicy(
|
||||
RestartClass.CLIENT_RECONNECT,
|
||||
"mcp.reconnect.client",
|
||||
BLAST_NONE,
|
||||
"none",
|
||||
False,
|
||||
"self_service",
|
||||
"record class, actor, client namespace, reason, and outcome",
|
||||
"Reconnect only the caller's client transport; no daemon or peer session changes.",
|
||||
ALL_REQUEST_ROLES,
|
||||
ALL_REQUEST_ROLES,
|
||||
),
|
||||
RestartClass.SESSION_RECONNECT: RestartClassPolicy(
|
||||
RestartClass.SESSION_RECONNECT,
|
||||
"mcp.reconnect.session",
|
||||
BLAST_LOW,
|
||||
"requesting_session_safe_point",
|
||||
False,
|
||||
"self_service",
|
||||
"record class, actor, session id, reason, and outcome",
|
||||
"Rebind identity, capability, and workspace state for one session.",
|
||||
ALL_REQUEST_ROLES,
|
||||
ALL_REQUEST_ROLES,
|
||||
),
|
||||
RestartClass.WORKER_RESTART: RestartClassPolicy(
|
||||
RestartClass.WORKER_RESTART,
|
||||
"mcp.restart.worker.request",
|
||||
BLAST_LOW,
|
||||
"target_worker",
|
||||
False,
|
||||
"controller_approval_and_automated_gates",
|
||||
"record class, actor, target worker, approval, drain proof, and outcome",
|
||||
"Restart one worker after its own lease and mutation scope is drained.",
|
||||
ALL_REQUEST_ROLES,
|
||||
("operator", "admin"),
|
||||
),
|
||||
RestartClass.ROLE_RUNTIME_RESTART: RestartClassPolicy(
|
||||
RestartClass.ROLE_RUNTIME_RESTART,
|
||||
"mcp.restart.role_runtime.request",
|
||||
BLAST_MEDIUM,
|
||||
"target_role_runtime",
|
||||
False,
|
||||
"controller_approval_and_automated_gates",
|
||||
"record class, actor, role namespace, approval, drain proof, and outcome",
|
||||
"Restart only the selected role runtime and then re-probe that namespace.",
|
||||
ALL_REQUEST_ROLES,
|
||||
("operator", "admin"),
|
||||
),
|
||||
RestartClass.CONNECTOR_RESTART: RestartClassPolicy(
|
||||
RestartClass.CONNECTOR_RESTART,
|
||||
"mcp.restart.connector.request",
|
||||
BLAST_MEDIUM,
|
||||
"target_connector",
|
||||
False,
|
||||
"controller_approval_and_automated_gates",
|
||||
"record class, actor, connector id, approval, drain proof, and outcome",
|
||||
"Restart one connector while unrelated role runtimes remain available.",
|
||||
ALL_REQUEST_ROLES,
|
||||
("operator", "admin"),
|
||||
),
|
||||
RestartClass.CONFIGURATION_RELOAD: RestartClassPolicy(
|
||||
RestartClass.CONFIGURATION_RELOAD,
|
||||
"mcp.reload.configuration.request",
|
||||
BLAST_LOW,
|
||||
"mutation_quiesce",
|
||||
False,
|
||||
"controller_approval_and_automated_gates",
|
||||
"record class, actor, configuration revision, approval, and outcome",
|
||||
"Gracefully reload configuration without replacing the daemon process.",
|
||||
ALL_REQUEST_ROLES,
|
||||
("operator", "admin"),
|
||||
),
|
||||
RestartClass.ROLLING_MCP_RESTART: RestartClassPolicy(
|
||||
RestartClass.ROLLING_MCP_RESTART,
|
||||
"mcp.restart.rolling.request",
|
||||
BLAST_MEDIUM,
|
||||
"one_instance_at_a_time",
|
||||
False,
|
||||
"controller_approval_and_automated_gates",
|
||||
"record class, actor, instance order, approval, per-instance drains, and outcome",
|
||||
"Drain, restart, verify, and restore one instance before advancing to the next.",
|
||||
CONTROL_ROLES,
|
||||
("operator", "admin"),
|
||||
),
|
||||
RestartClass.FULL_MCP_RESTART: RestartClassPolicy(
|
||||
RestartClass.FULL_MCP_RESTART,
|
||||
"mcp.restart.full.request",
|
||||
BLAST_HIGH,
|
||||
"all_sessions_and_mutations",
|
||||
True,
|
||||
"controller_approval_and_automated_gates",
|
||||
"record class, actor, full impact report, approval, drain proof, and outcome",
|
||||
"Stop and restore the complete MCP runtime only after a verified full drain.",
|
||||
CONTROL_ROLES,
|
||||
("operator", "admin"),
|
||||
),
|
||||
RestartClass.HOST_RESTART: RestartClassPolicy(
|
||||
RestartClass.HOST_RESTART,
|
||||
"mcp.restart.host.request",
|
||||
BLAST_HIGH,
|
||||
"all_host_work",
|
||||
True,
|
||||
"controller_approval_plus_infrastructure_operator",
|
||||
"record class, actor, host, incident or change id, approval, drain proof, and outcome",
|
||||
"Hand off to infrastructure ownership; reconcile every runtime after the host returns.",
|
||||
("controller", "operator", "admin"),
|
||||
("operator", "admin"),
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def resolve_restart_class(value: RestartClass | str) -> RestartClass:
|
||||
"""Resolve a restart class or fail closed for an unknown value."""
|
||||
|
||||
if isinstance(value, RestartClass):
|
||||
return value
|
||||
try:
|
||||
return RestartClass(str(value).strip())
|
||||
except ValueError as exc:
|
||||
raise ValueError(f"unknown restart class {value!r}; deny (fail closed)") from exc
|
||||
|
||||
|
||||
def restart_class_policy(value: RestartClass | str) -> RestartClassPolicy:
|
||||
"""Return the canonical policy for *value*."""
|
||||
|
||||
return RESTART_CLASS_POLICIES[resolve_restart_class(value)]
|
||||
|
||||
|
||||
def permissions_for_role(role: str | None) -> tuple[str, ...]:
|
||||
"""Return request permissions granted to a workflow role by this policy."""
|
||||
|
||||
normalized = str(role or "").strip().lower()
|
||||
return tuple(
|
||||
policy.required_permission
|
||||
for policy in RESTART_CLASS_POLICIES.values()
|
||||
if normalized in policy.request_roles
|
||||
)
|
||||
|
||||
|
||||
def _utc_now() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
@@ -75,6 +264,7 @@ class SessionImpact:
|
||||
heartbeat_stale: bool
|
||||
is_requester: bool
|
||||
live: bool
|
||||
connector: str | None = None
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
@@ -87,6 +277,7 @@ class SessionImpact:
|
||||
"heartbeat_stale": self.heartbeat_stale,
|
||||
"is_requester": self.is_requester,
|
||||
"live": self.live,
|
||||
"connector": self.connector,
|
||||
}
|
||||
|
||||
|
||||
@@ -105,6 +296,7 @@ class LeaseImpact:
|
||||
disruptive: bool
|
||||
is_mutation: bool
|
||||
is_critical_section: bool
|
||||
connector: str | None = None
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
@@ -119,6 +311,7 @@ class LeaseImpact:
|
||||
"disruptive": self.disruptive,
|
||||
"is_mutation": self.is_mutation,
|
||||
"is_critical_section": self.is_critical_section,
|
||||
"connector": self.connector,
|
||||
}
|
||||
|
||||
|
||||
@@ -127,6 +320,13 @@ class RestartImpactReport:
|
||||
"""Impact preview DTO returned to the console / operator (#642/#652)."""
|
||||
|
||||
coordinator_version: str
|
||||
restart_class: str
|
||||
restart_policy: dict[str, Any]
|
||||
policy_enforced: bool
|
||||
permission_authorized: bool
|
||||
role_authorized: bool
|
||||
approval_satisfied: bool
|
||||
authorization_reasons: list[str]
|
||||
evaluated_at: str
|
||||
dry_run: bool
|
||||
restart_performed: bool
|
||||
@@ -153,6 +353,13 @@ class RestartImpactReport:
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
"coordinator_version": self.coordinator_version,
|
||||
"restart_class": self.restart_class,
|
||||
"restart_policy": dict(self.restart_policy),
|
||||
"policy_enforced": self.policy_enforced,
|
||||
"permission_authorized": self.permission_authorized,
|
||||
"role_authorized": self.role_authorized,
|
||||
"approval_satisfied": self.approval_satisfied,
|
||||
"authorization_reasons": list(self.authorization_reasons),
|
||||
"evaluated_at": self.evaluated_at,
|
||||
"dry_run": self.dry_run,
|
||||
"restart_performed": self.restart_performed,
|
||||
@@ -206,6 +413,7 @@ def _classify_session(
|
||||
requesting_session_id and session_id == requesting_session_id
|
||||
),
|
||||
live=live,
|
||||
connector=(str(row.get("connector") or "").strip() or None),
|
||||
)
|
||||
|
||||
|
||||
@@ -258,6 +466,7 @@ def _classify_lease(row: Mapping[str, Any]) -> LeaseImpact:
|
||||
disruptive=disruptive,
|
||||
is_mutation=is_mutation,
|
||||
is_critical_section=disruptive,
|
||||
connector=(str(row.get("connector") or "").strip() or None),
|
||||
)
|
||||
|
||||
|
||||
@@ -279,6 +488,14 @@ def evaluate_restart_impact(
|
||||
requesting_session_id: str | None = None,
|
||||
dry_run: bool = True,
|
||||
session_heartbeat_stale_seconds: int = DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS,
|
||||
restart_class: RestartClass | str | None = None,
|
||||
requester_role: str | None = None,
|
||||
requester_permissions: Sequence[str] | None = None,
|
||||
controller_approved: bool = False,
|
||||
operator_authorized: bool = False,
|
||||
target_session_id: str | None = None,
|
||||
target_role: str | None = None,
|
||||
target_connector: str | None = None,
|
||||
) -> RestartImpactReport:
|
||||
"""Evaluate a proposed MCP restart and return an impact preview.
|
||||
|
||||
@@ -301,6 +518,61 @@ def evaluate_restart_impact(
|
||||
"""
|
||||
moment = now or _utc_now()
|
||||
reasons: list[str] = []
|
||||
authorization_reasons: list[str] = []
|
||||
|
||||
# ``None`` preserves the pre-#663 impact-only API for callers that have not
|
||||
# yet been migrated. All MCP requests pass an explicit class and therefore
|
||||
# take the fail-closed policy path.
|
||||
policy_enforced = restart_class is not None
|
||||
try:
|
||||
resolved_class = resolve_restart_class(
|
||||
restart_class or RestartClass.FULL_MCP_RESTART
|
||||
)
|
||||
policy = RESTART_CLASS_POLICIES[resolved_class]
|
||||
unknown_class = False
|
||||
except ValueError as exc:
|
||||
resolved_class = None
|
||||
policy = None
|
||||
unknown_class = True
|
||||
authorization_reasons.append(str(exc))
|
||||
|
||||
normalized_role = str(requester_role or "").strip().lower()
|
||||
granted = {str(p).strip() for p in (requester_permissions or ())}
|
||||
if policy_enforced and policy is not None:
|
||||
permission_authorized = policy.required_permission in granted
|
||||
role_authorized = normalized_role in policy.request_roles
|
||||
if not permission_authorized:
|
||||
authorization_reasons.append(
|
||||
f"missing required permission {policy.required_permission!r}"
|
||||
)
|
||||
if not role_authorized:
|
||||
authorization_reasons.append(
|
||||
f"role {normalized_role or 'unknown'!r} may not request "
|
||||
f"{policy.restart_class.value}"
|
||||
)
|
||||
elif unknown_class:
|
||||
permission_authorized = False
|
||||
role_authorized = False
|
||||
else:
|
||||
permission_authorized = True
|
||||
role_authorized = True
|
||||
|
||||
if policy_enforced and policy is not None:
|
||||
approval = policy.approval_requirement
|
||||
if approval == "self_service":
|
||||
approval_satisfied = True
|
||||
elif approval == "controller_approval_plus_infrastructure_operator":
|
||||
approval_satisfied = bool(controller_approved and operator_authorized)
|
||||
else:
|
||||
approval_satisfied = bool(controller_approved)
|
||||
if not approval_satisfied:
|
||||
authorization_reasons.append(
|
||||
f"approval requirement not satisfied: {approval}"
|
||||
)
|
||||
elif unknown_class:
|
||||
approval_satisfied = False
|
||||
else:
|
||||
approval_satisfied = True
|
||||
|
||||
inventory_complete = bool(inventory.get("inventory_complete", False))
|
||||
incomplete_reasons = [str(r) for r in (inventory.get("incomplete_reasons") or [])]
|
||||
@@ -323,15 +595,67 @@ def evaluate_restart_impact(
|
||||
]
|
||||
lease_impacts = [_classify_lease(l) for l in leases_raw]
|
||||
|
||||
# Only *other* live sessions and live leases constitute blast radius: a
|
||||
# restart that would kill only the requesting session with no other work in
|
||||
# flight is safe.
|
||||
# Route impact through the selected class. Narrow classes never inherit a
|
||||
# full-runtime drain merely because unrelated work exists.
|
||||
target_complete = True
|
||||
if resolved_class in {
|
||||
RestartClass.CLIENT_RECONNECT,
|
||||
RestartClass.SESSION_RECONNECT,
|
||||
RestartClass.CONFIGURATION_RELOAD,
|
||||
}:
|
||||
scoped_sessions: list[SessionImpact] = []
|
||||
scoped_leases: list[LeaseImpact] = []
|
||||
elif resolved_class == RestartClass.WORKER_RESTART:
|
||||
selected_session = (target_session_id or "").strip()
|
||||
target_complete = bool(selected_session)
|
||||
scoped_sessions = [
|
||||
s for s in session_impacts if s.session_id == selected_session
|
||||
]
|
||||
scoped_leases = [
|
||||
l for l in lease_impacts if l.session_id == selected_session
|
||||
]
|
||||
elif resolved_class == RestartClass.ROLE_RUNTIME_RESTART:
|
||||
selected_role = (target_role or "").strip().lower()
|
||||
target_complete = bool(selected_role)
|
||||
scoped_sessions = [
|
||||
s for s in session_impacts if str(s.role or "").lower() == selected_role
|
||||
]
|
||||
scoped_leases = [
|
||||
l for l in lease_impacts if str(l.role or "").lower() == selected_role
|
||||
]
|
||||
elif resolved_class == RestartClass.CONNECTOR_RESTART:
|
||||
selected_connector = (target_connector or "").strip()
|
||||
target_complete = bool(selected_connector)
|
||||
scoped_sessions = [
|
||||
s for s in session_impacts if s.connector == selected_connector
|
||||
]
|
||||
scoped_leases = [
|
||||
l for l in lease_impacts if l.connector == selected_connector
|
||||
]
|
||||
else:
|
||||
scoped_sessions = list(session_impacts)
|
||||
scoped_leases = list(lease_impacts)
|
||||
|
||||
if policy_enforced and not target_complete:
|
||||
authorization_reasons.append(
|
||||
f"target required for {resolved_class.value if resolved_class else 'unknown class'}"
|
||||
)
|
||||
|
||||
other_live_sessions = [
|
||||
s for s in session_impacts if s.live and not s.is_requester
|
||||
s for s in scoped_sessions if s.live and not s.is_requester
|
||||
]
|
||||
disruptive_leases = [l for l in lease_impacts if l.disruptive]
|
||||
critical_sections = [l for l in lease_impacts if l.is_critical_section]
|
||||
mutations = [l for l in lease_impacts if l.is_mutation]
|
||||
disruptive_leases = [l for l in scoped_leases if l.disruptive]
|
||||
critical_sections = [l for l in scoped_leases if l.is_critical_section]
|
||||
mutations = [l for l in scoped_leases if l.is_mutation]
|
||||
terminal_lock_in_scope = (
|
||||
terminal_lock
|
||||
if resolved_class
|
||||
not in {
|
||||
RestartClass.CLIENT_RECONNECT,
|
||||
RestartClass.SESSION_RECONNECT,
|
||||
}
|
||||
else None
|
||||
)
|
||||
|
||||
affected_issues = sorted(
|
||||
{
|
||||
@@ -348,9 +672,24 @@ def evaluate_restart_impact(
|
||||
}
|
||||
)
|
||||
|
||||
disruptive = bool(disruptive_leases or other_live_sessions or terminal_lock)
|
||||
disruptive = bool(
|
||||
disruptive_leases or other_live_sessions or terminal_lock_in_scope
|
||||
)
|
||||
|
||||
if not inventory_complete:
|
||||
authorization_ok = bool(
|
||||
not unknown_class
|
||||
and permission_authorized
|
||||
and role_authorized
|
||||
and approval_satisfied
|
||||
and target_complete
|
||||
)
|
||||
|
||||
if policy_enforced and not authorization_ok:
|
||||
verdict = VERDICT_UNSAFE
|
||||
allow_restart = False
|
||||
reasons.append("restart class authorization denied (fail closed)")
|
||||
reasons.extend(authorization_reasons)
|
||||
elif not inventory_complete:
|
||||
verdict = VERDICT_UNSAFE
|
||||
allow_restart = False
|
||||
reasons.append(
|
||||
@@ -381,7 +720,7 @@ def evaluate_restart_impact(
|
||||
f"{len(critical_sections)} critical section(s) in flight "
|
||||
"(active lease with a live owner)"
|
||||
)
|
||||
if terminal_lock:
|
||||
if terminal_lock_in_scope:
|
||||
reasons.append("active terminal (merge) lock present")
|
||||
|
||||
override_would_allow = bool(inventory_complete and disruptive)
|
||||
@@ -411,6 +750,12 @@ def evaluate_restart_impact(
|
||||
audit_record = {
|
||||
"event": "restart_impact_evaluated",
|
||||
"coordinator_version": COORDINATOR_VERSION,
|
||||
"restart_class": (
|
||||
resolved_class.value if resolved_class else str(restart_class or "")
|
||||
),
|
||||
"required_permission": (
|
||||
policy.required_permission if policy is not None else None
|
||||
),
|
||||
"evaluated_at": moment.isoformat(),
|
||||
"dry_run": dry_run,
|
||||
"operator_override": bool(operator_override),
|
||||
@@ -424,6 +769,15 @@ def evaluate_restart_impact(
|
||||
|
||||
return RestartImpactReport(
|
||||
coordinator_version=COORDINATOR_VERSION,
|
||||
restart_class=(
|
||||
resolved_class.value if resolved_class else str(restart_class or "")
|
||||
),
|
||||
restart_policy=policy.as_dict() if policy is not None else {},
|
||||
policy_enforced=policy_enforced,
|
||||
permission_authorized=permission_authorized,
|
||||
role_authorized=role_authorized,
|
||||
approval_satisfied=approval_satisfied,
|
||||
authorization_reasons=authorization_reasons,
|
||||
evaluated_at=moment.isoformat(),
|
||||
dry_run=dry_run,
|
||||
restart_performed=False,
|
||||
@@ -440,9 +794,11 @@ def evaluate_restart_impact(
|
||||
affected_issues=affected_issues,
|
||||
affected_prs=affected_prs,
|
||||
mutations=mutations,
|
||||
terminal_lock=dict(terminal_lock)
|
||||
if isinstance(terminal_lock, Mapping)
|
||||
else terminal_lock,
|
||||
terminal_lock=(
|
||||
dict(terminal_lock_in_scope)
|
||||
if isinstance(terminal_lock_in_scope, Mapping)
|
||||
else terminal_lock_in_scope
|
||||
),
|
||||
ack_state=ack_state,
|
||||
prior_recovery_attempts=prior_recovery_attempts,
|
||||
counts=counts,
|
||||
|
||||
@@ -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()
|
||||
@@ -0,0 +1,232 @@
|
||||
"""Permission, drain, routing, and audit matrix for restart classes (#663)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from datetime import datetime, timezone
|
||||
|
||||
import restart_coordinator as rc
|
||||
|
||||
NOW = datetime(2026, 7, 24, 20, 0, tzinfo=timezone.utc)
|
||||
|
||||
|
||||
def _inventory() -> dict:
|
||||
return {
|
||||
"inventory_complete": True,
|
||||
"sessions": [
|
||||
{
|
||||
"session_id": "requester",
|
||||
"role": "author",
|
||||
"profile": "prgs-author",
|
||||
"pid": os.getpid(),
|
||||
"status": "active",
|
||||
"last_heartbeat_at": NOW.isoformat(),
|
||||
},
|
||||
{
|
||||
"session_id": "reviewer",
|
||||
"role": "reviewer",
|
||||
"profile": "prgs-reviewer",
|
||||
"pid": os.getpid(),
|
||||
"status": "active",
|
||||
"last_heartbeat_at": NOW.isoformat(),
|
||||
},
|
||||
],
|
||||
"leases": [
|
||||
{
|
||||
"lease_id": "review-lease",
|
||||
"session_id": "reviewer",
|
||||
"role": "reviewer",
|
||||
"phase": "reviewing",
|
||||
"work_kind": "pr",
|
||||
"work_number": 900,
|
||||
"worktree_path": "/tmp/review-900",
|
||||
"freshness": {"freshness": "active"},
|
||||
}
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
def _evaluate(
|
||||
restart_class: rc.RestartClass,
|
||||
*,
|
||||
role: str = "controller",
|
||||
permissions: tuple[str, ...] | None = None,
|
||||
approved: bool = True,
|
||||
operator: bool = True,
|
||||
**targets,
|
||||
):
|
||||
return rc.evaluate_restart_impact(
|
||||
_inventory(),
|
||||
now=NOW,
|
||||
requesting_session_id="requester",
|
||||
restart_class=restart_class,
|
||||
requester_role=role,
|
||||
requester_permissions=(
|
||||
permissions if permissions is not None
|
||||
else rc.permissions_for_role(role)
|
||||
),
|
||||
controller_approved=approved,
|
||||
operator_authorized=operator,
|
||||
**targets,
|
||||
)
|
||||
|
||||
|
||||
def test_policy_table_covers_exactly_all_nine_classes():
|
||||
assert set(rc.RESTART_CLASS_POLICIES) == set(rc.RestartClass)
|
||||
assert len(rc.RESTART_CLASS_POLICIES) == 9
|
||||
for restart_class, policy in rc.RESTART_CLASS_POLICIES.items():
|
||||
assert policy.restart_class is restart_class
|
||||
assert policy.required_permission
|
||||
assert policy.expected_blast_radius in {
|
||||
rc.BLAST_NONE, rc.BLAST_LOW, rc.BLAST_MEDIUM, rc.BLAST_HIGH
|
||||
}
|
||||
assert policy.drain_requirement
|
||||
assert policy.approval_requirement
|
||||
assert policy.audit_requirement
|
||||
assert policy.recovery_behavior
|
||||
|
||||
|
||||
def test_permission_matrix_allows_each_class_with_exact_permission():
|
||||
targets = {
|
||||
rc.RestartClass.WORKER_RESTART: {"target_session_id": "reviewer"},
|
||||
rc.RestartClass.ROLE_RUNTIME_RESTART: {"target_role": "reviewer"},
|
||||
rc.RestartClass.CONNECTOR_RESTART: {"target_connector": "github"},
|
||||
}
|
||||
for restart_class, policy in rc.RESTART_CLASS_POLICIES.items():
|
||||
report = _evaluate(
|
||||
restart_class,
|
||||
permissions=(policy.required_permission,),
|
||||
**targets.get(restart_class, {}),
|
||||
)
|
||||
assert report.permission_authorized, restart_class
|
||||
assert report.role_authorized, restart_class
|
||||
assert report.approval_satisfied, restart_class
|
||||
assert report.audit_record["restart_class"] == restart_class.value
|
||||
assert (
|
||||
report.audit_record["required_permission"]
|
||||
== policy.required_permission
|
||||
)
|
||||
|
||||
|
||||
def test_missing_or_nearby_permission_denies():
|
||||
report = _evaluate(
|
||||
rc.RestartClass.ROLE_RUNTIME_RESTART,
|
||||
permissions=("mcp.restart.worker.request",),
|
||||
target_role="reviewer",
|
||||
)
|
||||
assert report.verdict == rc.VERDICT_UNSAFE
|
||||
assert not report.allow_restart
|
||||
assert not report.permission_authorized
|
||||
assert any("missing required permission" in r for r in report.reasons)
|
||||
|
||||
|
||||
def test_unknown_restart_class_denies_fail_closed():
|
||||
report = rc.evaluate_restart_impact(
|
||||
_inventory(),
|
||||
now=NOW,
|
||||
restart_class="surprise_reboot",
|
||||
requester_role="admin",
|
||||
requester_permissions=("mcp.restart.host.request",),
|
||||
controller_approved=True,
|
||||
operator_authorized=True,
|
||||
)
|
||||
assert report.verdict == rc.VERDICT_UNSAFE
|
||||
assert not report.allow_restart
|
||||
assert report.restart_policy == {}
|
||||
assert any("unknown restart class" in r for r in report.reasons)
|
||||
|
||||
|
||||
def test_worker_roles_cannot_request_full_or_host_restart():
|
||||
for role in rc.WORKER_ROLES:
|
||||
granted = rc.permissions_for_role(role)
|
||||
assert "mcp.restart.full.request" not in granted
|
||||
assert "mcp.restart.host.request" not in granted
|
||||
report = _evaluate(
|
||||
rc.RestartClass.FULL_MCP_RESTART,
|
||||
role=role,
|
||||
permissions=granted,
|
||||
)
|
||||
assert not report.role_authorized
|
||||
assert not report.allow_restart
|
||||
|
||||
|
||||
def test_controller_approval_is_independent_of_permission():
|
||||
report = _evaluate(
|
||||
rc.RestartClass.WORKER_RESTART,
|
||||
approved=False,
|
||||
target_session_id="reviewer",
|
||||
)
|
||||
assert report.permission_authorized
|
||||
assert not report.approval_satisfied
|
||||
assert not report.allow_restart
|
||||
|
||||
|
||||
def test_narrow_classes_do_not_inherit_full_drain_or_peer_lease_block():
|
||||
for restart_class in (
|
||||
rc.RestartClass.CLIENT_RECONNECT,
|
||||
rc.RestartClass.SESSION_RECONNECT,
|
||||
rc.RestartClass.CONFIGURATION_RELOAD,
|
||||
):
|
||||
report = _evaluate(restart_class)
|
||||
assert not report.restart_policy["full_drain_required"]
|
||||
assert report.counts["leases_disruptive"] == 0
|
||||
assert report.counts["sessions_live_other"] == 0
|
||||
assert report.counts["critical_sections"] == 0
|
||||
assert report.counts["mutations"] == 0
|
||||
assert report.allow_restart, (restart_class, report.reasons)
|
||||
|
||||
|
||||
def test_client_reconnect_does_not_wait_for_unrelated_terminal_lock():
|
||||
inventory = _inventory()
|
||||
inventory["terminal_lock"] = {"terminal_pr": 901}
|
||||
report = rc.evaluate_restart_impact(
|
||||
inventory,
|
||||
now=NOW,
|
||||
requesting_session_id="requester",
|
||||
restart_class=rc.RestartClass.CLIENT_RECONNECT,
|
||||
requester_role="author",
|
||||
requester_permissions=rc.permissions_for_role("author"),
|
||||
)
|
||||
assert report.allow_restart
|
||||
assert report.terminal_lock is None
|
||||
|
||||
|
||||
def test_scoped_restart_only_counts_named_target():
|
||||
report = _evaluate(
|
||||
rc.RestartClass.ROLE_RUNTIME_RESTART,
|
||||
target_role="author",
|
||||
)
|
||||
assert report.counts["leases_disruptive"] == 0
|
||||
assert report.affected_prs == []
|
||||
assert report.allow_restart
|
||||
|
||||
reviewer = _evaluate(
|
||||
rc.RestartClass.ROLE_RUNTIME_RESTART,
|
||||
target_role="reviewer",
|
||||
)
|
||||
assert reviewer.counts["leases_disruptive"] == 1
|
||||
assert reviewer.affected_prs == [900]
|
||||
assert not reviewer.allow_restart
|
||||
|
||||
|
||||
def test_missing_scoped_target_denies_instead_of_widening():
|
||||
for restart_class in (
|
||||
rc.RestartClass.WORKER_RESTART,
|
||||
rc.RestartClass.ROLE_RUNTIME_RESTART,
|
||||
rc.RestartClass.CONNECTOR_RESTART,
|
||||
):
|
||||
report = _evaluate(restart_class)
|
||||
assert not report.allow_restart
|
||||
assert any("target required" in r for r in report.reasons)
|
||||
|
||||
|
||||
def test_only_full_and_host_classes_require_full_drain():
|
||||
requiring_full = {
|
||||
restart_class
|
||||
for restart_class, policy in rc.RESTART_CLASS_POLICIES.items()
|
||||
if policy.full_drain_required
|
||||
}
|
||||
assert requiring_full == {
|
||||
rc.RestartClass.FULL_MCP_RESTART,
|
||||
rc.RestartClass.HOST_RESTART,
|
||||
}
|
||||
@@ -444,6 +444,12 @@ class TestAuditEmission(unittest.TestCase):
|
||||
)
|
||||
self.assertEqual(record["target"]["namespace"], NAMESPACE)
|
||||
self.assertEqual(record["target"]["mode"], "restart")
|
||||
self.assertEqual(
|
||||
record["target"]["restart_class"], "role_runtime_restart"
|
||||
)
|
||||
self.assertEqual(
|
||||
record["metadata"]["restart_class"], "role_runtime_restart"
|
||||
)
|
||||
self.assertEqual(record["result"], console_audit.RESULT_ALLOWED)
|
||||
self.assertEqual(record["actor"]["subject"], "[email protected]")
|
||||
self.assertFalse(record["metadata"]["process_kill_executed"])
|
||||
|
||||
@@ -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 B1–B5).
|
||||
"""
|
||||
|
||||
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()
|
||||
@@ -42,6 +42,8 @@ from webui.lease_loader import load_lease_snapshot, snapshot_to_dict as lease_sn
|
||||
from webui.lease_views import render_leases_page
|
||||
from webui.queue_loader import load_queue_snapshot, snapshot_to_dict as queue_snapshot_to_dict
|
||||
from webui.queue_views import render_queue_page
|
||||
from webui.traffic_loader import load_traffic_snapshot, snapshot_to_dict as traffic_snapshot_to_dict
|
||||
from webui.traffic_views import render_traffic_page
|
||||
from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as worktree_snapshot_to_dict
|
||||
from webui.worktree_views import render_worktrees_page
|
||||
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
|
||||
@@ -80,6 +82,7 @@ def _stub_page(title: str, description: str) -> HTMLResponse:
|
||||
|
||||
|
||||
_LEGACY_PAGES = (
|
||||
("/traffic", "Traffic", "workflow traffic-control view (#640)"),
|
||||
("/queue", "Queue", "live PR and issue dashboard (#429)"),
|
||||
("/projects", "Projects", "registry and onboarding (#427)"),
|
||||
("/prompts", "Prompts", "canonical workflow prompt library (#428)"),
|
||||
@@ -200,6 +203,15 @@ async def api_queue(_request: Request) -> JSONResponse:
|
||||
return JSONResponse(queue_snapshot_to_dict(load_queue_snapshot()))
|
||||
|
||||
|
||||
async def traffic(_request: Request) -> HTMLResponse:
|
||||
snapshot = load_traffic_snapshot()
|
||||
return HTMLResponse(render_traffic_page(snapshot))
|
||||
|
||||
|
||||
async def api_traffic(_request: Request) -> JSONResponse:
|
||||
return JSONResponse(traffic_snapshot_to_dict(load_traffic_snapshot()))
|
||||
|
||||
|
||||
def _load_project_registry() -> tuple[ProjectRegistry | None, RegistryError | None]:
|
||||
"""Load the registry, converting validation failure into a fail-closed pair."""
|
||||
try:
|
||||
@@ -736,6 +748,8 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
|
||||
Route("/system-health", system_health, methods=["GET"]),
|
||||
Route("/queue", queue, methods=["GET"]),
|
||||
Route("/api/queue", api_queue, methods=["GET"]),
|
||||
Route("/traffic", traffic, methods=["GET"]),
|
||||
Route("/api/traffic", api_traffic, methods=["GET"]),
|
||||
Route("/projects", projects, methods=["GET"]),
|
||||
Route("/projects/{project_id}", project_detail, methods=["GET"]),
|
||||
Route("/api/projects", api_projects, methods=["GET"]),
|
||||
|
||||
@@ -182,10 +182,17 @@ def _extract_reviewer_leases(
|
||||
parsed = parse_reviewer_lease_comment(comment.get("body") or "")
|
||||
if not parsed:
|
||||
continue
|
||||
subject_pr = parsed.get("pr_number") or pr_number
|
||||
leases.append(
|
||||
{
|
||||
**parsed,
|
||||
"pr_number": parsed.get("pr_number") or pr_number,
|
||||
"pr_number": subject_pr,
|
||||
# The lease subject is the PR, never the linked issue: a
|
||||
# reviewer lease on PR #N must not be attributed to issue #N
|
||||
# or to the issue that PR closes (#640).
|
||||
"kind": "pr",
|
||||
"number": subject_pr,
|
||||
"role": "reviewer",
|
||||
"comment_id": comment.get("id"),
|
||||
"author": (comment.get("user") or {}).get("login"),
|
||||
"created_at": comment.get("created_at"),
|
||||
|
||||
@@ -41,6 +41,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
|
||||
NavItem("/system-health", "System health"),
|
||||
)),
|
||||
NavGroup("Traffic", (
|
||||
NavItem("/traffic", "Traffic control"),
|
||||
NavItem("/queue", "Queue"),
|
||||
NavItem("/leases", "Leases"),
|
||||
NavItem("/actions", "Actions"),
|
||||
|
||||
+34
-3
@@ -4,7 +4,7 @@ from __future__ import annotations
|
||||
|
||||
import os
|
||||
import re
|
||||
from dataclasses import dataclass
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Callable
|
||||
from urllib.parse import urlparse
|
||||
@@ -31,10 +31,20 @@ class PaginationMeta:
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class QueueItem:
|
||||
"""One queue row.
|
||||
|
||||
``extra`` holds *display* strings for the queue page (values are truncated
|
||||
or humanized for rendering). ``signals`` holds the *authoritative* typed
|
||||
values taken straight from the Gitea payload, for consumers that classify
|
||||
or pin state rather than render it (#640). Never derive identity or
|
||||
concurrency decisions from ``extra``.
|
||||
"""
|
||||
|
||||
number: int
|
||||
title: str
|
||||
badges: tuple[str, ...]
|
||||
extra: dict[str, str]
|
||||
signals: dict[str, Any] = field(default_factory=dict)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
@@ -134,21 +144,37 @@ def _format_pr_item(pr: dict, badges: tuple[str, ...]) -> QueueItem:
|
||||
"mergeable" if mergeable is True else "conflicted" if mergeable is False else "unknown"
|
||||
)
|
||||
linked = _extract_linked_issue(pr.get("title"), pr.get("body"))
|
||||
head_sha = str(head.get("sha") or "")
|
||||
labels = tuple(
|
||||
str(lb.get("name") or "") for lb in (pr.get("labels") or []) if lb.get("name")
|
||||
)
|
||||
return QueueItem(
|
||||
number=int(pr["number"]),
|
||||
title=str(pr.get("title") or ""),
|
||||
badges=badges,
|
||||
extra={
|
||||
"branch": f"{head.get('ref', '?')} → {base.get('ref', '?')}",
|
||||
"head_sha": str(head.get("sha") or "")[:12],
|
||||
# Display only — truncated. Pin against signals["head_sha"] instead.
|
||||
"head_sha": head_sha[:12],
|
||||
"mergeable": merge_label,
|
||||
"linked_issue": str(linked) if linked is not None else "",
|
||||
},
|
||||
signals={
|
||||
"head_sha": head_sha,
|
||||
"head_ref": str(head.get("ref") or ""),
|
||||
"base_ref": str(base.get("ref") or ""),
|
||||
"mergeable": mergeable if isinstance(mergeable, bool) else None,
|
||||
"labels": labels,
|
||||
"linked_issue": linked,
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
|
||||
labels = ", ".join(lb.get("name", "") for lb in issue.get("labels", []))
|
||||
label_names = tuple(
|
||||
str(lb.get("name") or "") for lb in (issue.get("labels") or []) if lb.get("name")
|
||||
)
|
||||
labels = ", ".join(label_names)
|
||||
assignee = (issue.get("assignee") or {}).get("login", "")
|
||||
return QueueItem(
|
||||
number=int(issue["number"]),
|
||||
@@ -159,6 +185,11 @@ def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
|
||||
"assignee": assignee or "unassigned",
|
||||
"state": str(issue.get("state") or ""),
|
||||
},
|
||||
signals={
|
||||
"labels": label_names,
|
||||
"assignee": assignee,
|
||||
"state": str(issue.get("state") or ""),
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -38,6 +38,7 @@ from dataclasses import asdict, dataclass
|
||||
from typing import Any
|
||||
|
||||
import mcp_namespace_health
|
||||
import restart_coordinator
|
||||
import runtime_recovery_guard
|
||||
from webui import console_audit, console_authz
|
||||
|
||||
@@ -99,6 +100,14 @@ def _clean(value: Any) -> str:
|
||||
return str(value or "").strip()
|
||||
|
||||
|
||||
def restart_class_for_mode(mode: str) -> str:
|
||||
"""Map the existing namespace controls onto the #663 class taxonomy."""
|
||||
|
||||
if _clean(mode) == MODE_RELOAD:
|
||||
return restart_coordinator.RestartClass.CONFIGURATION_RELOAD.value
|
||||
return restart_coordinator.RestartClass.ROLE_RUNTIME_RESTART.value
|
||||
|
||||
|
||||
# --- Mutation ledger --------------------------------------------------------
|
||||
|
||||
|
||||
@@ -256,6 +265,7 @@ def build_restart_preview(
|
||||
|
||||
return {
|
||||
"action_id": action_id,
|
||||
"restart_class": restart_class_for_mode(md),
|
||||
"namespace": ns,
|
||||
"mode": md,
|
||||
"scope_valid": scope_error is None,
|
||||
@@ -309,6 +319,7 @@ def assess_restart_request(
|
||||
"reason_code": reason_code,
|
||||
"detail": detail,
|
||||
"action_id": action_id,
|
||||
"restart_class": restart_class_for_mode(md),
|
||||
"namespace": ns,
|
||||
"mode": md,
|
||||
"preview": preview,
|
||||
@@ -393,6 +404,7 @@ def assess_restart_request(
|
||||
"process."
|
||||
),
|
||||
"action_id": action_id,
|
||||
"restart_class": restart_class_for_mode(md),
|
||||
"namespace": ns,
|
||||
"mode": md,
|
||||
"preview": preview,
|
||||
@@ -441,7 +453,11 @@ def execute_restart(
|
||||
else console_audit.RESULT_DENIED
|
||||
),
|
||||
principal=principal,
|
||||
target={"namespace": assessment["namespace"], "mode": assessment["mode"]},
|
||||
target={
|
||||
"namespace": assessment["namespace"],
|
||||
"mode": assessment["mode"],
|
||||
"restart_class": assessment["restart_class"],
|
||||
},
|
||||
reason_code=assessment["reason_code"],
|
||||
detail=assessment["detail"],
|
||||
request_id=request_id,
|
||||
@@ -450,6 +466,7 @@ def execute_restart(
|
||||
"gates_passed": assessment["gates_passed"],
|
||||
"process_kill_executed": False,
|
||||
"post_restart_verification_required": True,
|
||||
"restart_class": assessment["restart_class"],
|
||||
},
|
||||
)
|
||||
|
||||
@@ -463,6 +480,7 @@ def execute_restart(
|
||||
"namespace": assessment["namespace"],
|
||||
"mode": assessment["mode"],
|
||||
"action_id": action_id,
|
||||
"restart_class": assessment["restart_class"],
|
||||
"process_kill_executed": False,
|
||||
"host_hook": assessment["preview"]["restart_hook"],
|
||||
"next_action": (
|
||||
|
||||
@@ -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()
|
||||
@@ -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)
|
||||
Reference in New Issue
Block a user