Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6e6ca94338 | ||
|
|
2f4dec8323 | ||
|
|
3a9d634c17 | ||
|
|
2068bae341 | ||
|
|
7af40fb5ff | ||
|
|
9517834913 | ||
|
|
824c42f7e3 | ||
|
|
578c44b685 | ||
|
|
3b68d15593 | ||
|
|
a4c73766f4 | ||
|
|
9f686253eb | ||
|
|
b2e28428a4 | ||
|
|
95e4aae287 | ||
|
|
dac40ab9b3 | ||
|
|
ccde9e8f11 | ||
|
|
1948d3dc21 | ||
|
|
069a9af7e6 | ||
|
|
1cbbde0089 | ||
|
|
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
|
# Body phrases that prove an issue is an implementation container, not a
|
||||||
# unit of direct author work (#844). Matched case-insensitively against the
|
# unit of direct author work (#844 / #854). Matched case-insensitively against
|
||||||
# issue body. Title alone is never sufficient (ordinary issues may mention
|
# the issue body. Title alone is never sufficient (ordinary issues may mention
|
||||||
# "epic" incidentally).
|
# "epic", "roadmap", "vision", or "umbrella" incidentally).
|
||||||
|
#
|
||||||
|
# #854 extends the #844 marker set so product-vision (#652), phased-roadmap
|
||||||
|
# (#653), and umbrella (#655) coordination records — which do not use the word
|
||||||
|
# "epic" — are classified with the same semantic exclusion as epic containers.
|
||||||
_CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
|
_CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
|
||||||
|
# Epic / child-only (#844, live #631)
|
||||||
"implementation is delivered via child issues only",
|
"implementation is delivered via child issues only",
|
||||||
"implementation is delivered through child issues only",
|
"implementation is delivered through child issues only",
|
||||||
"implementation is delivered via child issues",
|
"implementation is delivered via child issues",
|
||||||
@@ -150,9 +155,26 @@ _CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
|
|||||||
"coordination container",
|
"coordination container",
|
||||||
"child-only container",
|
"child-only container",
|
||||||
"implementation is delegated to child",
|
"implementation is delegated to child",
|
||||||
|
# Vision / roadmap / umbrella coordination (#854, live #652/#653/#655).
|
||||||
|
# Prefer authoritative non-implementation / child-only scope language over
|
||||||
|
# bare words like "roadmap" so ordinary implementable issues that mention
|
||||||
|
# a parent vision or roadmap stay eligible.
|
||||||
|
"do not implement features on this issue",
|
||||||
|
"implementing features on this roadmap issue",
|
||||||
|
"implementation is via linked children only",
|
||||||
|
"no product feature claimed complete on this issue alone",
|
||||||
|
"phased delivery roadmap and epic sequencing",
|
||||||
|
"this issue is the enduring source of truth",
|
||||||
|
"enduring source of truth for the",
|
||||||
|
"canonical product vision — enduring source of truth",
|
||||||
|
"canonical product vision - enduring source of truth",
|
||||||
|
"state: vision-active",
|
||||||
|
"state: roadmap-active",
|
||||||
)
|
)
|
||||||
|
|
||||||
# Explicit epic / umbrella labels (structured evidence preferred over title).
|
# Explicit epic / umbrella / vision / roadmap labels (structured evidence
|
||||||
|
# preferred over title). Tracker alone is *not* included — ordinary issues
|
||||||
|
# may carry a tracker label without being non-implementable containers.
|
||||||
_EPIC_LABELS: frozenset[str] = frozenset(
|
_EPIC_LABELS: frozenset[str] = frozenset(
|
||||||
{
|
{
|
||||||
"type:epic",
|
"type:epic",
|
||||||
@@ -161,6 +183,16 @@ _EPIC_LABELS: frozenset[str] = frozenset(
|
|||||||
"scope:epic",
|
"scope:epic",
|
||||||
"type:umbrella",
|
"type:umbrella",
|
||||||
"umbrella",
|
"umbrella",
|
||||||
|
"kind:umbrella",
|
||||||
|
"scope:umbrella",
|
||||||
|
"type:vision",
|
||||||
|
"vision",
|
||||||
|
"kind:vision",
|
||||||
|
"scope:vision",
|
||||||
|
"type:roadmap",
|
||||||
|
"roadmap",
|
||||||
|
"kind:roadmap",
|
||||||
|
"scope:roadmap",
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -222,16 +254,45 @@ class WorkCandidate:
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _title_container_prefix(title_l: str) -> str | None:
|
||||||
|
"""Return a coordination-title prefix token if *title_l* uses one (#854).
|
||||||
|
|
||||||
|
Title prefixes alone never exclude; they only corroborate body/label
|
||||||
|
evidence. Ordinary issues may say "roadmap" or "vision" mid-title.
|
||||||
|
"""
|
||||||
|
for prefix, token in (
|
||||||
|
("epic:", "title_epic_prefix"),
|
||||||
|
("epic ", "title_epic_prefix"),
|
||||||
|
("umbrella:", "title_umbrella_prefix"),
|
||||||
|
("umbrella ", "title_umbrella_prefix"),
|
||||||
|
("roadmap:", "title_roadmap_prefix"),
|
||||||
|
("roadmap ", "title_roadmap_prefix"),
|
||||||
|
("product vision:", "title_vision_prefix"),
|
||||||
|
("product vision ", "title_vision_prefix"),
|
||||||
|
("vision:", "title_vision_prefix"),
|
||||||
|
("vision ", "title_vision_prefix"),
|
||||||
|
):
|
||||||
|
if title_l.startswith(prefix):
|
||||||
|
return token
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
def classify_epic_or_child_only_container(
|
def classify_epic_or_child_only_container(
|
||||||
c: WorkCandidate,
|
c: WorkCandidate,
|
||||||
) -> tuple[bool, str | None]:
|
) -> tuple[bool, str | None]:
|
||||||
"""Return whether *c* is an epic / child-only implementation container (#844).
|
"""Return whether *c* is a non-implementable coordination container (#844/#854).
|
||||||
|
|
||||||
Exclusion uses structured evidence first (labels, body scope language).
|
Exclusion uses structured evidence first (labels, body scope language).
|
||||||
A bare title containing the word "epic" is **not** enough — ordinary
|
A bare title containing the words "epic", "roadmap", "vision", or
|
||||||
implementable issues may mention epics incidentally. A title that is
|
"umbrella" is **not** enough — ordinary implementable issues may mention
|
||||||
explicitly prefixed ``Epic:`` only counts when the body also proves
|
those terms incidentally. Explicit title prefixes (``Epic:``, ``Roadmap:``,
|
||||||
child-only / no-direct-implementation scope (or an epic label is present).
|
``Product vision:``, ``Umbrella:``) only count when the body also proves
|
||||||
|
child-only / no-direct-implementation scope (or a container label is
|
||||||
|
present).
|
||||||
|
|
||||||
|
Covers epic, product-vision, phased-roadmap, umbrella, and child-only
|
||||||
|
records so the allocator never assigns coordination containers as direct
|
||||||
|
author work.
|
||||||
|
|
||||||
PRs are never classified as containers here (they already have a head).
|
PRs are never classified as containers here (they already have a head).
|
||||||
"""
|
"""
|
||||||
@@ -245,25 +306,28 @@ def classify_epic_or_child_only_container(
|
|||||||
title_l = title.lower()
|
title_l = title.lower()
|
||||||
|
|
||||||
body_hits = [m for m in _CHILD_ONLY_BODY_MARKERS if m in body_l]
|
body_hits = [m for m in _CHILD_ONLY_BODY_MARKERS if m in body_l]
|
||||||
title_epic_prefix = title_l.startswith("epic:") or title_l.startswith("epic ")
|
title_prefix = _title_container_prefix(title_l)
|
||||||
|
|
||||||
if epic_label:
|
if epic_label:
|
||||||
detail = f"label={epic_label[0]}"
|
detail = f"label={epic_label[0]}"
|
||||||
if body_hits:
|
if body_hits:
|
||||||
detail = f"{detail}; body_marker={body_hits[0]!r}"
|
detail = f"{detail}; body_marker={body_hits[0]!r}"
|
||||||
|
if title_prefix:
|
||||||
|
detail = f"{title_prefix}; {detail}"
|
||||||
return True, detail
|
return True, detail
|
||||||
|
|
||||||
if body_hits:
|
if body_hits:
|
||||||
# Body proves child-only / umbrella scope. Title "Epic:" is corroborating
|
# Body proves child-only / vision / roadmap / umbrella scope. Title
|
||||||
# but not required — containers without the word still exclude.
|
# prefixes are corroborating but not required — containers without the
|
||||||
|
# title word still exclude.
|
||||||
detail = f"body_marker={body_hits[0]!r}"
|
detail = f"body_marker={body_hits[0]!r}"
|
||||||
if title_epic_prefix:
|
if title_prefix:
|
||||||
detail = f"title_epic_prefix; {detail}"
|
detail = f"{title_prefix}; {detail}"
|
||||||
return True, detail
|
return True, detail
|
||||||
|
|
||||||
# Title-only "Epic:" without body scope evidence is insufficient (#844 AC:
|
# Title-only coordination prefix without body scope evidence is
|
||||||
# eligibility does not rely solely on the word "Epic" in a title).
|
# insufficient (#844/#854 AC: eligibility does not rely solely on a title
|
||||||
# Similarly, incidental "epic" mid-title without markers stays eligible.
|
# word). Incidental mid-title mentions without markers stay eligible.
|
||||||
return False, None
|
return False, None
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -57,6 +57,8 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
|
|||||||
| `/system-health` | System-health dashboard — readiness, version/uptime, dependencies, MCP namespaces, stale-runtime parity (#639) |
|
| `/system-health` | System-health dashboard — readiness, version/uptime, dependencies, MCP namespaces, stale-runtime parity (#639) |
|
||||||
| `/queue` | Live PR and issue queue dashboard (#429) |
|
| `/queue` | Live PR and issue queue dashboard (#429) |
|
||||||
| `/api/queue` | JSON queue export with pagination metadata |
|
| `/api/queue` | JSON queue export with pagination metadata |
|
||||||
|
| `/traffic` | Workflow traffic-control view — runnable, leased, blocked, needs-controller, terminal-complete (#640) |
|
||||||
|
| `/api/traffic` | JSON traffic-control export with state classifications and next safe role actions |
|
||||||
| `/projects` | Project registry list with status and onboarding progress (#427, #635) |
|
| `/projects` | Project registry list with status and onboarding progress (#427, #635) |
|
||||||
| `/projects/{id}` | Project detail + onboarding checklist |
|
| `/projects/{id}` | Project detail + onboarding checklist |
|
||||||
| `/api/v1/projects` | Versioned JSON registry export (#635) |
|
| `/api/v1/projects` | Versioned JSON registry export (#635) |
|
||||||
@@ -85,6 +87,35 @@ Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with
|
|||||||
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
|
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
|
||||||
local validator preview only (no Gitea mutations, no server-side storage).
|
local validator preview only (no Gitea mutations, no server-side storage).
|
||||||
|
|
||||||
|
### Traffic-control state vocabulary (#640)
|
||||||
|
|
||||||
|
The traffic view classifies each open issue/PR into exactly one bucket:
|
||||||
|
|
||||||
|
| Bucket | Meaning | Operator implication |
|
||||||
|
|--------|---------|----------------------|
|
||||||
|
| **runnable** | No active lease, no block reason, safe for its expected role | Next role may start work |
|
||||||
|
| **leased** | Active author claim or reviewer PR lease | Do not stomp; wait or adopt via role tools |
|
||||||
|
| **blocked** | Dependency, missing head pin, conflict, or unmet dependency | Author remediation first |
|
||||||
|
| **needs_controller** | Contaminated, controller-only diagnosis, or `status:blocked` | Controller only |
|
||||||
|
| **terminal_complete** | Reconciler / terminal-lock territory | Reconciler cleanup path |
|
||||||
|
|
||||||
|
`status:blocked` items route to **needs_controller**, not **blocked**:
|
||||||
|
`expected_role_for_candidate` sends them to the controller, and the blocker
|
||||||
|
reason renders in either bucket.
|
||||||
|
|
||||||
|
**Live path contracts (do not invent):**
|
||||||
|
|
||||||
|
- PR head pins come from `QueueItem.signals["head_sha"]` (full SHA). Display
|
||||||
|
`extra["head_sha"]` is truncated and must never be used for routing.
|
||||||
|
- Reviewer leases are keyed as `(pr, pr_number)` only — never via a linked
|
||||||
|
`issue_number` on the same lease marker.
|
||||||
|
- Issue claims come from `claim_inventory["entries"]`
|
||||||
|
(`issue_claim_heartbeat.build_claim_inventory`). There is no `active_claims`
|
||||||
|
key.
|
||||||
|
- Queue display badges are only: `blocked`, `claimed`, `duplicate`, `stale`,
|
||||||
|
`in-review`, `open`. Review verdicts (`request-changes`, `approved`) are
|
||||||
|
**not** queue badges; traffic does not invent them from the queue loader.
|
||||||
|
|
||||||
## System health API (#634)
|
## System health API (#634)
|
||||||
|
|
||||||
`GET /api/v1/system/health` is the structured, read-only health surface for
|
`GET /api/v1/system/health` is the structured, read-only health surface for
|
||||||
|
|||||||
+1015
File diff suppressed because it is too large
Load Diff
+354
-78
@@ -2068,6 +2068,7 @@ import lease_lifecycle # noqa: E402
|
|||||||
import lease_policy # noqa: E402
|
import lease_policy # noqa: E402
|
||||||
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
|
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
|
||||||
import restart_coordinator # noqa: E402 # #658 MCP restart coordinator/impact
|
import restart_coordinator # noqa: E402 # #658 MCP restart coordinator/impact
|
||||||
|
import drain_proof # noqa: E402 # #661 pre-restart drain proof and hard gate
|
||||||
import incident_bridge # noqa: E402
|
import incident_bridge # noqa: E402
|
||||||
import sentry_observability # noqa: E402 (#606 optional Sentry observability)
|
import sentry_observability # noqa: E402 (#606 optional Sentry observability)
|
||||||
import sentry_incident_bridge # noqa: E402 (#607 Sentry→Gitea incident bridge)
|
import sentry_incident_bridge # noqa: E402 (#607 Sentry→Gitea incident bridge)
|
||||||
@@ -9551,15 +9552,13 @@ def gitea_edit_pr(
|
|||||||
if closing:
|
if closing:
|
||||||
gate_reasons = _profile_operation_gate("gitea.pr.close")
|
gate_reasons = _profile_operation_gate("gitea.pr.close")
|
||||||
if gate_reasons:
|
if gate_reasons:
|
||||||
return {
|
return _build_operation_gate_refusal(
|
||||||
"success": False,
|
"gitea.pr.close",
|
||||||
"performed": False,
|
gate_reasons,
|
||||||
"pr_number": pr_number,
|
pr_number=pr_number,
|
||||||
"requested_state": "closed",
|
requested_state="closed",
|
||||||
"required_permission": "gitea.pr.close",
|
required_permission="gitea.pr.close",
|
||||||
"reasons": gate_reasons,
|
)
|
||||||
"permission_report": _permission_block_report("gitea.pr.close"),
|
|
||||||
}
|
|
||||||
|
|
||||||
h, o, r = _resolve(remote, host, org, repo)
|
h, o, r = _resolve(remote, host, org, repo)
|
||||||
auth = _auth(h)
|
auth = _auth(h)
|
||||||
@@ -13822,13 +13821,17 @@ def gitea_view_issue(
|
|||||||
|
|
||||||
def _permission_block_report(required_operation: str,
|
def _permission_block_report(required_operation: str,
|
||||||
identity: str | None = None) -> dict:
|
identity: str | None = None) -> dict:
|
||||||
"""Structured, LLM-safe explanation of a permission denial (#142).
|
"""Structured, LLM-safe explanation of a permission denial (#142, #897).
|
||||||
|
|
||||||
Built only after a gate has already refused; it adds guidance to the
|
Built only after a gate has already refused; it adds guidance to the
|
||||||
refusal and never widens any permission, performs network I/O, or
|
refusal and never widens any permission, performs network I/O, or
|
||||||
raises (fail-soft: degrades to a minimal fail-closed report). Names
|
raises (fail-soft: degrades to a minimal fail-closed report). Names
|
||||||
configured profiles only — never auth references, tokens, endpoint
|
configured profiles only — never auth references, tokens, endpoint
|
||||||
URLs, or keychain IDs.
|
URLs, or keychain IDs.
|
||||||
|
|
||||||
|
#897: never fabricate a missing permission when the active profile
|
||||||
|
already allows the operation. That path is a diagnostic defect (the
|
||||||
|
refusal was not a permission denial), not a cue to switch profiles.
|
||||||
"""
|
"""
|
||||||
report = {
|
report = {
|
||||||
"requested_operation": required_operation,
|
"requested_operation": required_operation,
|
||||||
@@ -13840,6 +13843,7 @@ def _permission_block_report(required_operation: str,
|
|||||||
"matching_configured_profiles": [],
|
"matching_configured_profiles": [],
|
||||||
"runtime_switching_supported": False,
|
"runtime_switching_supported": False,
|
||||||
"different_mcp_namespace_required": True,
|
"different_mcp_namespace_required": True,
|
||||||
|
"diagnostic_defect": False,
|
||||||
"exact_safe_next_action": (
|
"exact_safe_next_action": (
|
||||||
"Ask the operator to fix GITEA_MCP_CONFIG/GITEA_MCP_PROFILE; "
|
"Ask the operator to fix GITEA_MCP_CONFIG/GITEA_MCP_PROFILE; "
|
||||||
"the active profile could not be resolved (fail closed)."),
|
"the active profile could not be resolved (fail closed)."),
|
||||||
@@ -13854,6 +13858,32 @@ def _permission_block_report(required_operation: str,
|
|||||||
report["active_allowed_operations"] = (
|
report["active_allowed_operations"] = (
|
||||||
profile.get("allowed_operations") or [])
|
profile.get("allowed_operations") or [])
|
||||||
|
|
||||||
|
# #897: fail closed as a diagnostic defect when the active profile
|
||||||
|
# already holds the operation — callers must not invent "missing".
|
||||||
|
try:
|
||||||
|
holds, _hold_reason = gitea_config.check_operation(
|
||||||
|
required_operation,
|
||||||
|
profile.get("allowed_operations") or [],
|
||||||
|
profile.get("forbidden_operations") or [],
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
holds = False
|
||||||
|
if holds:
|
||||||
|
report["missing_permission"] = None
|
||||||
|
report["required_permission"] = required_operation
|
||||||
|
report["diagnostic_defect"] = True
|
||||||
|
report["different_mcp_namespace_required"] = False
|
||||||
|
report["exact_safe_next_action"] = (
|
||||||
|
"Diagnostic defect: the active profile already allows "
|
||||||
|
f"{required_operation}. This is not a permission denial — "
|
||||||
|
"inspect blocker_kind / reasons (stale-runtime or runtime-mode). "
|
||||||
|
"Do not call gitea_activate_profile or switch MCP sessions."
|
||||||
|
)
|
||||||
|
report["matching_configured_profiles"] = [
|
||||||
|
p for p in [profile.get("profile_name")] if p
|
||||||
|
]
|
||||||
|
return report
|
||||||
|
|
||||||
matching = []
|
matching = []
|
||||||
try:
|
try:
|
||||||
config = gitea_config.load_config() or {}
|
config = gitea_config.load_config() or {}
|
||||||
@@ -13902,6 +13932,205 @@ def _permission_block_report(required_operation: str,
|
|||||||
return report
|
return report
|
||||||
|
|
||||||
|
|
||||||
|
def _reason_is_stale_runtime(reason: str) -> bool:
|
||||||
|
"""True when *reason* is a master-parity / stale-daemon refusal (#897)."""
|
||||||
|
r = (reason or "").lower()
|
||||||
|
if not r:
|
||||||
|
return False
|
||||||
|
if "stale relative to live master" in r:
|
||||||
|
return True
|
||||||
|
if "server code is stale" in r:
|
||||||
|
return True
|
||||||
|
if "daemon is stale" in r:
|
||||||
|
return True
|
||||||
|
if "started at commit" in r and "workspace master is now" in r:
|
||||||
|
return True
|
||||||
|
if "mcp server started at" in r and "stale" in r:
|
||||||
|
return True
|
||||||
|
if "restart the server to load the current capability gates" in r:
|
||||||
|
return True
|
||||||
|
if "restart/reconnect before mutating" in r:
|
||||||
|
return True
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def _reason_is_runtime_mode(reason: str) -> bool:
|
||||||
|
"""True when *reason* is a stable-control / runtime-mode refusal (#897)."""
|
||||||
|
r = (reason or "").lower()
|
||||||
|
if not r:
|
||||||
|
return False
|
||||||
|
if _reason_is_stale_runtime(reason):
|
||||||
|
return False
|
||||||
|
if "runtime mode could not be assessed" in r:
|
||||||
|
return True
|
||||||
|
if "runtime mode is" in r:
|
||||||
|
return True
|
||||||
|
if "stable control runtime" in r:
|
||||||
|
return True
|
||||||
|
if "dev-test" in r and ("runtime" in r or "production" in r):
|
||||||
|
return True
|
||||||
|
if "development worktree" in r or "dev worktree" in r:
|
||||||
|
return True
|
||||||
|
if "launched from a 'branches/" in r or "launched from a \"branches/" in r:
|
||||||
|
return True
|
||||||
|
if "process-root / active-workspace alignment" in r:
|
||||||
|
return True
|
||||||
|
if "namespace" in r and "reproof" in r:
|
||||||
|
return True
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def _reason_is_permission(reason: str) -> bool:
|
||||||
|
"""True when *reason* is a genuine profile-permission denial (#897)."""
|
||||||
|
r = (reason or "").lower()
|
||||||
|
if not r:
|
||||||
|
return False
|
||||||
|
if _reason_is_stale_runtime(reason) or _reason_is_runtime_mode(reason):
|
||||||
|
return False
|
||||||
|
if "profile could not be resolved" in r:
|
||||||
|
return True
|
||||||
|
if "profile has no configured allowed operations" in r:
|
||||||
|
return True
|
||||||
|
if "profile forbids" in r:
|
||||||
|
return True
|
||||||
|
if "profile is not allowed to" in r:
|
||||||
|
return True
|
||||||
|
if "unrecognized forbidden operation" in r:
|
||||||
|
return True
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def _classify_operation_gate_reasons(reasons: list[str]) -> dict:
|
||||||
|
"""Partition gate reasons into stale / runtime-mode / permission (#897)."""
|
||||||
|
stale: list[str] = []
|
||||||
|
runtime_mode: list[str] = []
|
||||||
|
permission: list[str] = []
|
||||||
|
other: list[str] = []
|
||||||
|
for reason in reasons or []:
|
||||||
|
if _reason_is_stale_runtime(reason):
|
||||||
|
stale.append(reason)
|
||||||
|
elif _reason_is_runtime_mode(reason):
|
||||||
|
runtime_mode.append(reason)
|
||||||
|
elif _reason_is_permission(reason):
|
||||||
|
permission.append(reason)
|
||||||
|
else:
|
||||||
|
other.append(reason)
|
||||||
|
return {
|
||||||
|
"stale_runtime": stale,
|
||||||
|
"runtime_mode": runtime_mode,
|
||||||
|
"permission": permission,
|
||||||
|
"other": other,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _stale_runtime_reconnect_action() -> str:
|
||||||
|
"""Sanctioned recovery for a stale daemon — reconnect only (#685/#897)."""
|
||||||
|
return (
|
||||||
|
"Reconnect the IDE/client MCP session so the server reloads at the "
|
||||||
|
"current master head. Do not call gitea_activate_profile or switch "
|
||||||
|
"MCP role sessions — profile switching does not clear a stale daemon."
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _build_operation_gate_refusal(
|
||||||
|
required_operation: str,
|
||||||
|
reasons: list[str],
|
||||||
|
**extra_fields,
|
||||||
|
) -> dict:
|
||||||
|
"""Structured gate refusal with typed blockers (#897).
|
||||||
|
|
||||||
|
Stale-runtime and runtime-mode refusals never attach a
|
||||||
|
``permission_report`` and never recommend profile switching. True
|
||||||
|
permission denials still get ``permission_report``. When both apply,
|
||||||
|
causes are reported separately under distinct fields.
|
||||||
|
"""
|
||||||
|
classified = _classify_operation_gate_reasons(reasons)
|
||||||
|
stale = classified["stale_runtime"]
|
||||||
|
runtime_mode = classified["runtime_mode"]
|
||||||
|
permission = classified["permission"]
|
||||||
|
other = classified["other"]
|
||||||
|
|
||||||
|
blocked: dict = {
|
||||||
|
"success": False,
|
||||||
|
"performed": False,
|
||||||
|
"reasons": list(reasons),
|
||||||
|
"mutation_performed": False,
|
||||||
|
"session_context_audit": session_ctx.mutation_context_audit_fields(),
|
||||||
|
"gate_reason_classes": {
|
||||||
|
"stale_runtime": list(stale),
|
||||||
|
"runtime_mode": list(runtime_mode),
|
||||||
|
"permission": list(permission),
|
||||||
|
"other": list(other),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
if stale:
|
||||||
|
parity = _current_master_parity()
|
||||||
|
blocked["blocker_kind"] = "runtime_reconnect_required"
|
||||||
|
blocked["restart_required"] = True
|
||||||
|
blocked["stop_required"] = True
|
||||||
|
blocked["startup_head"] = parity.get("startup_head")
|
||||||
|
blocked["current_head"] = parity.get("current_head")
|
||||||
|
blocked["daemon_start_head"] = (
|
||||||
|
parity.get("daemon_start_head") or parity.get("startup_head")
|
||||||
|
)
|
||||||
|
blocked["local_head"] = (
|
||||||
|
parity.get("local_head") or parity.get("current_head")
|
||||||
|
)
|
||||||
|
blocked["live_remote_head"] = parity.get("live_remote_head")
|
||||||
|
blocked["live_stale"] = bool(parity.get("live_stale"))
|
||||||
|
blocked["live_known"] = bool(parity.get("live_known"))
|
||||||
|
blocked["exact_safe_next_action"] = _stale_runtime_reconnect_action()
|
||||||
|
if permission or other:
|
||||||
|
blocked["permission_block_reasons"] = list(permission) + list(other)
|
||||||
|
blocked["stale_runtime_reasons"] = list(stale)
|
||||||
|
# Never attach permission_report for a staleness refusal.
|
||||||
|
blocked.update(extra_fields)
|
||||||
|
return blocked
|
||||||
|
|
||||||
|
if runtime_mode:
|
||||||
|
blocked["blocker_kind"] = "runtime_mode_blocked"
|
||||||
|
blocked["restart_required"] = False
|
||||||
|
blocked["stop_required"] = True
|
||||||
|
blocked["exact_safe_next_action"] = (
|
||||||
|
"Real workflow mutations run only on the promoted stable control "
|
||||||
|
"runtime. Promote/reload the stable runtime; do not call "
|
||||||
|
"gitea_activate_profile or switch MCP role sessions to clear a "
|
||||||
|
"runtime-mode block."
|
||||||
|
)
|
||||||
|
if permission or other:
|
||||||
|
blocked["permission_block_reasons"] = list(permission) + list(other)
|
||||||
|
blocked["runtime_mode_reasons"] = list(runtime_mode)
|
||||||
|
blocked.update(extra_fields)
|
||||||
|
return blocked
|
||||||
|
|
||||||
|
# Pure permission (or unclassified-as-permission) denial.
|
||||||
|
blocked["blocker_kind"] = "permission_denied"
|
||||||
|
blocked["permission_report"] = _permission_block_report(required_operation)
|
||||||
|
blocked.update(extra_fields)
|
||||||
|
return blocked
|
||||||
|
|
||||||
|
|
||||||
|
def _permission_report_for_gate_reasons(
|
||||||
|
required_operation: str,
|
||||||
|
reasons: list[str] | None,
|
||||||
|
) -> dict | None:
|
||||||
|
"""Attach ``permission_report`` only for true permission denials (#897).
|
||||||
|
|
||||||
|
Call sites that historically always attached a permission report after
|
||||||
|
``_profile_operation_gate`` should use this so stale/runtime refusals
|
||||||
|
do not emit a fabricated missing-permission payload.
|
||||||
|
"""
|
||||||
|
if not reasons:
|
||||||
|
return None
|
||||||
|
classified = _classify_operation_gate_reasons(reasons)
|
||||||
|
if classified["stale_runtime"] or classified["runtime_mode"]:
|
||||||
|
return None
|
||||||
|
if not (classified["permission"] or classified["other"]):
|
||||||
|
return None
|
||||||
|
return _permission_block_report(required_operation)
|
||||||
|
|
||||||
|
|
||||||
def _role_for_operation(op: str) -> str | None:
|
def _role_for_operation(op: str) -> str | None:
|
||||||
# Normalize op first
|
# Normalize op first
|
||||||
try:
|
try:
|
||||||
@@ -14078,7 +14307,7 @@ def _master_parity_block(op: str) -> list[str]:
|
|||||||
|
|
||||||
|
|
||||||
def _profile_operation_gate(op: str) -> list[str]:
|
def _profile_operation_gate(op: str) -> list[str]:
|
||||||
"""Profile permission check for a single gated operation (#126, #216, #420).
|
"""Profile permission check for a single gated operation (#126, #216, #420, #897).
|
||||||
|
|
||||||
Issue discussion comments are gated separately from the gitea.pr.*
|
Issue discussion comments are gated separately from the gitea.pr.*
|
||||||
review/merge family: listing requires ``gitea.read``, creating requires
|
review/merge family: listing requires ``gitea.read``, creating requires
|
||||||
@@ -14091,21 +14320,26 @@ def _profile_operation_gate(op: str) -> list[str]:
|
|||||||
capability gate that has since been merged, and when the runtime itself is
|
capability gate that has since been merged, and when the runtime itself is
|
||||||
not the promoted stable control runtime (#615) -- a dev/test or unknown
|
not the promoted stable control runtime (#615) -- a dev/test or unknown
|
||||||
runtime holds production credentials but has not been promoted.
|
runtime holds production credentials but has not been promoted.
|
||||||
|
|
||||||
|
#897: collect *all* independent refusal classes (stale, runtime-mode,
|
||||||
|
permission) rather than short-circuiting after the first. Callers that
|
||||||
|
only need a boolean still treat any non-empty list as blocked; typed
|
||||||
|
consumers (``_build_operation_gate_refusal``) can separate causes.
|
||||||
"""
|
"""
|
||||||
stale_reasons = _master_parity_block(op)
|
reasons: list[str] = []
|
||||||
if stale_reasons:
|
reasons.extend(_master_parity_block(op))
|
||||||
return stale_reasons
|
reasons.extend(_runtime_mode_block(op))
|
||||||
runtime_reasons = _runtime_mode_block(op)
|
|
||||||
if runtime_reasons:
|
|
||||||
return runtime_reasons
|
|
||||||
try:
|
try:
|
||||||
profile = get_profile()
|
profile = get_profile()
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
return [f"profile could not be resolved (fail closed): {_redact(str(exc))}"]
|
reasons.append(
|
||||||
|
f"profile could not be resolved (fail closed): {_redact(str(exc))}"
|
||||||
|
)
|
||||||
|
return reasons
|
||||||
op_ok, op_reason = gitea_config.check_operation(
|
op_ok, op_reason = gitea_config.check_operation(
|
||||||
op, profile["allowed_operations"], profile["forbidden_operations"])
|
op, profile["allowed_operations"], profile["forbidden_operations"])
|
||||||
if op_ok:
|
if op_ok:
|
||||||
return []
|
return reasons
|
||||||
|
|
||||||
if _try_auto_switch_for_operation(op):
|
if _try_auto_switch_for_operation(op):
|
||||||
try:
|
try:
|
||||||
@@ -14113,17 +14347,26 @@ def _profile_operation_gate(op: str) -> list[str]:
|
|||||||
op_ok, op_reason = gitea_config.check_operation(
|
op_ok, op_reason = gitea_config.check_operation(
|
||||||
op, profile["allowed_operations"], profile["forbidden_operations"])
|
op, profile["allowed_operations"], profile["forbidden_operations"])
|
||||||
if op_ok:
|
if op_ok:
|
||||||
return []
|
return reasons
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
return [f"profile could not be resolved (fail closed): {_redact(str(exc))}"]
|
reasons.append(
|
||||||
|
f"profile could not be resolved (fail closed): {_redact(str(exc))}"
|
||||||
|
)
|
||||||
|
return reasons
|
||||||
|
|
||||||
if op_reason == "no-allowed-operations":
|
if op_reason == "no-allowed-operations":
|
||||||
return ["profile has no configured allowed operations (fail closed)"]
|
reasons.append(
|
||||||
if op_reason == "forbidden":
|
"profile has no configured allowed operations (fail closed)"
|
||||||
return [f"profile forbids '{op}'"]
|
)
|
||||||
if op_reason == "invalid-forbidden-entry":
|
elif op_reason == "forbidden":
|
||||||
return ["profile has an unrecognized forbidden operation entry (fail closed)"]
|
reasons.append(f"profile forbids '{op}'")
|
||||||
return [f"profile is not allowed to {op}"]
|
elif op_reason == "invalid-forbidden-entry":
|
||||||
|
reasons.append(
|
||||||
|
"profile has an unrecognized forbidden operation entry (fail closed)"
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
reasons.append(f"profile is not allowed to {op}")
|
||||||
|
return reasons
|
||||||
|
|
||||||
|
|
||||||
def _mutation_config_authority_block(required_operation: str) -> dict | None:
|
def _mutation_config_authority_block(required_operation: str) -> dict | None:
|
||||||
@@ -14343,10 +14586,14 @@ def _session_context_mutation_block(
|
|||||||
|
|
||||||
|
|
||||||
def _profile_permission_block(required_operation: str, **extra_fields) -> dict | None:
|
def _profile_permission_block(required_operation: str, **extra_fields) -> dict | None:
|
||||||
"""Structured permission denial for gated tools (#69, #142).
|
"""Structured operation-gate denial for gated tools (#69, #142, #897).
|
||||||
|
|
||||||
Returns a block dict when the active profile forbids *required_operation*,
|
Returns a block dict when the active profile forbids *required_operation*,
|
||||||
or ``None`` when the gate passes. Never performs network I/O.
|
the daemon is stale, or the runtime mode is not mutation-safe — or
|
||||||
|
``None`` when the gate passes. Never performs network I/O.
|
||||||
|
|
||||||
|
#897: stale-runtime and runtime-mode refusals are typed
|
||||||
|
(``blocker_kind``) and never carry a ``permission_report``.
|
||||||
"""
|
"""
|
||||||
req_role = "reviewer" if any(required_operation.startswith(p) for p in (
|
req_role = "reviewer" if any(required_operation.startswith(p) for p in (
|
||||||
"gitea.pr.approve", "gitea.pr.merge", "gitea.pr.request_changes", "gitea.pr.review"
|
"gitea.pr.approve", "gitea.pr.merge", "gitea.pr.request_changes", "gitea.pr.review"
|
||||||
@@ -14356,15 +14603,9 @@ def _profile_permission_block(required_operation: str, **extra_fields) -> dict |
|
|||||||
|
|
||||||
reasons = _profile_operation_gate(required_operation)
|
reasons = _profile_operation_gate(required_operation)
|
||||||
if reasons:
|
if reasons:
|
||||||
blocked = {
|
return _build_operation_gate_refusal(
|
||||||
"success": False,
|
required_operation, reasons, **extra_fields
|
||||||
"performed": False,
|
)
|
||||||
"reasons": reasons,
|
|
||||||
"permission_report": _permission_block_report(required_operation),
|
|
||||||
"session_context_audit": session_ctx.mutation_context_audit_fields(),
|
|
||||||
}
|
|
||||||
blocked.update(extra_fields)
|
|
||||||
return blocked
|
|
||||||
|
|
||||||
auth_block = _mutation_config_authority_block(required_operation)
|
auth_block = _mutation_config_authority_block(required_operation)
|
||||||
if auth_block is not None:
|
if auth_block is not None:
|
||||||
@@ -14493,20 +14734,14 @@ def gitea_acquire_reviewer_pr_lease(
|
|||||||
"""Acquire a per-PR reviewer lease before review/merge mutations (#407)."""
|
"""Acquire a per-PR reviewer lease before review/merge mutations (#407)."""
|
||||||
read_block = _profile_operation_gate("gitea.read")
|
read_block = _profile_operation_gate("gitea.read")
|
||||||
if read_block:
|
if read_block:
|
||||||
return {
|
return _build_operation_gate_refusal(
|
||||||
"success": False,
|
"gitea.read", read_block, acquired=False
|
||||||
"acquired": False,
|
)
|
||||||
"reasons": read_block,
|
|
||||||
"permission_report": _permission_block_report("gitea.read"),
|
|
||||||
}
|
|
||||||
comment_block = _profile_operation_gate("gitea.pr.comment")
|
comment_block = _profile_operation_gate("gitea.pr.comment")
|
||||||
if comment_block:
|
if comment_block:
|
||||||
return {
|
return _build_operation_gate_refusal(
|
||||||
"success": False,
|
"gitea.pr.comment", comment_block, acquired=False
|
||||||
"acquired": False,
|
)
|
||||||
"reasons": comment_block,
|
|
||||||
"permission_report": _permission_block_report("gitea.pr.comment"),
|
|
||||||
}
|
|
||||||
|
|
||||||
# task=acquire_reviewer_pr_lease so verify_preflight_purity runs shared #604
|
# task=acquire_reviewer_pr_lease so verify_preflight_purity runs shared #604
|
||||||
# anti-stomp for the declared lease-acquire mutation inventory entry.
|
# anti-stomp for the declared lease-acquire mutation inventory entry.
|
||||||
@@ -14622,20 +14857,14 @@ def gitea_acquire_merger_pr_lease(
|
|||||||
"""
|
"""
|
||||||
read_block = _profile_operation_gate("gitea.read")
|
read_block = _profile_operation_gate("gitea.read")
|
||||||
if read_block:
|
if read_block:
|
||||||
return {
|
return _build_operation_gate_refusal(
|
||||||
"success": False,
|
"gitea.read", read_block, acquired=False
|
||||||
"acquired": False,
|
)
|
||||||
"reasons": read_block,
|
|
||||||
"permission_report": _permission_block_report("gitea.read"),
|
|
||||||
}
|
|
||||||
comment_block = _profile_operation_gate("gitea.pr.comment")
|
comment_block = _profile_operation_gate("gitea.pr.comment")
|
||||||
if comment_block:
|
if comment_block:
|
||||||
return {
|
return _build_operation_gate_refusal(
|
||||||
"success": False,
|
"gitea.pr.comment", comment_block, acquired=False
|
||||||
"acquired": False,
|
)
|
||||||
"reasons": comment_block,
|
|
||||||
"permission_report": _permission_block_report("gitea.pr.comment"),
|
|
||||||
}
|
|
||||||
merge_block = _profile_operation_gate("gitea.pr.merge")
|
merge_block = _profile_operation_gate("gitea.pr.merge")
|
||||||
if merge_block:
|
if merge_block:
|
||||||
return {
|
return {
|
||||||
@@ -19374,14 +19603,12 @@ def gitea_update_pr_branch_by_merge(
|
|||||||
# Permission: author branch push / PR mutation surface.
|
# Permission: author branch push / PR mutation surface.
|
||||||
push_block = _profile_operation_gate("gitea.branch.push")
|
push_block = _profile_operation_gate("gitea.branch.push")
|
||||||
if push_block:
|
if push_block:
|
||||||
return {
|
return _build_operation_gate_refusal(
|
||||||
"success": False,
|
"gitea.branch.push",
|
||||||
"performed": False,
|
push_block,
|
||||||
"mutation_allowed": False,
|
mutation_allowed=False,
|
||||||
"reasons": push_block,
|
role_kind=role,
|
||||||
"permission_report": _permission_block_report("gitea.branch.push"),
|
)
|
||||||
"role_kind": role,
|
|
||||||
}
|
|
||||||
|
|
||||||
if role != "author":
|
if role != "author":
|
||||||
pre = pr_sync_status.assess_update_pr_branch_preflight(
|
pre = pr_sync_status.assess_update_pr_branch_preflight(
|
||||||
@@ -22342,6 +22569,8 @@ def gitea_request_mcp_restart(
|
|||||||
request_override: bool = False,
|
request_override: bool = False,
|
||||||
session_id: str | None = None,
|
session_id: str | None = None,
|
||||||
limit: int = 200,
|
limit: int = 200,
|
||||||
|
drain_proof_json: str | None = None,
|
||||||
|
request_break_glass: bool = False,
|
||||||
) -> dict:
|
) -> dict:
|
||||||
"""Evaluate a proposed MCP restart and return an impact preview (#658).
|
"""Evaluate a proposed MCP restart and return an impact preview (#658).
|
||||||
|
|
||||||
@@ -22351,10 +22580,15 @@ def gitea_request_mcp_restart(
|
|||||||
verdict, so the console (#642/#652) and operators can see what a restart
|
verdict, so the console (#642/#652) and operators can see what a restart
|
||||||
would disrupt *before* any concurrent LLM work is destroyed.
|
would disrupt *before* any concurrent LLM work is destroyed.
|
||||||
|
|
||||||
This tool is **dry-run and never restarts anything.** The mutative apply
|
This tool **never restarts a process.** In dry-run (the default) it returns
|
||||||
path is a separate child gated by a drain proof (non-goal here); calling
|
only the impact preview. With ``dry_run=False`` it enforces the #661 hard
|
||||||
with ``dry_run=False`` still performs no restart and reports that apply is
|
gate: the apply request must present a valid, unexpired, clean drain proof
|
||||||
not yet available.
|
(``drain_proof_json``) or it is denied and a durable incident descriptor is
|
||||||
|
returned under ``incident``. Break-glass is the only bypass and is honoured
|
||||||
|
only when ``request_break_glass`` is set *and* the environment carries
|
||||||
|
``GITEA_BREAKGLASS_RESTART_AUTHORIZATION``. Even an authorized gate performs
|
||||||
|
no restart here; actual execution is a further child. The gate outcome is
|
||||||
|
reported under ``apply_gate`` / ``apply_authorized``.
|
||||||
|
|
||||||
Operator override authority is read from the process environment
|
Operator override authority is read from the process environment
|
||||||
(``GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION``), never self-asserted by
|
(``GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION``), never self-asserted by
|
||||||
@@ -22473,12 +22707,54 @@ def gitea_request_mcp_restart(
|
|||||||
payload["requesting_session_id"] = sid
|
payload["requesting_session_id"] = sid
|
||||||
payload["operator_override_requested"] = bool(request_override)
|
payload["operator_override_requested"] = bool(request_override)
|
||||||
payload["operator_override_authorized"] = operator_authorized
|
payload["operator_override_authorized"] = operator_authorized
|
||||||
|
# Actual restart execution remains a further child; this tool never restarts
|
||||||
|
# a process. What #661 adds is the *hard gate*: an apply request (dry_run
|
||||||
|
# False) must present a valid, unexpired, clean drain proof, or it is denied
|
||||||
|
# and a durable incident is raised. Break-glass is the only bypass and its
|
||||||
|
# authorization is read from the environment, never self-asserted.
|
||||||
payload["apply_supported"] = False
|
payload["apply_supported"] = False
|
||||||
if not dry_run:
|
if not dry_run:
|
||||||
payload["reasons"] = list(payload.get("reasons") or []) + [
|
proof_obj: dict | None = None
|
||||||
"apply requested but not supported: sanctioned restart apply is "
|
proof_parse_error: str | None = None
|
||||||
"gated by a drain proof (separate child); no restart performed (#658)"
|
if drain_proof_json:
|
||||||
]
|
try:
|
||||||
|
parsed = json.loads(drain_proof_json)
|
||||||
|
proof_obj = parsed if isinstance(parsed, dict) else None
|
||||||
|
if proof_obj is None:
|
||||||
|
proof_parse_error = "drain_proof_json is not a JSON object"
|
||||||
|
except (ValueError, TypeError) as exc:
|
||||||
|
proof_parse_error = f"invalid drain_proof_json: {_redact(str(exc))}"
|
||||||
|
|
||||||
|
break_glass_authorized = bool(
|
||||||
|
(
|
||||||
|
os.environ.get("GITEA_BREAKGLASS_RESTART_AUTHORIZATION") or ""
|
||||||
|
).strip()
|
||||||
|
)
|
||||||
|
break_glass = bool(request_break_glass and break_glass_authorized)
|
||||||
|
|
||||||
|
expected_fp = drain_proof.impact_fingerprint(report.as_dict())
|
||||||
|
gate = drain_proof.gate_apply_restart(
|
||||||
|
proof=proof_obj,
|
||||||
|
break_glass=break_glass,
|
||||||
|
expected_impact_fingerprint=expected_fp,
|
||||||
|
requesting_session_id=sid,
|
||||||
|
)
|
||||||
|
gate_payload = gate.as_dict()
|
||||||
|
if proof_parse_error and not break_glass:
|
||||||
|
gate_payload["reasons"] = [proof_parse_error] + list(
|
||||||
|
gate_payload.get("reasons") or []
|
||||||
|
)
|
||||||
|
payload["apply_gate"] = gate_payload
|
||||||
|
payload["apply_authorized"] = gate.allow
|
||||||
|
payload["break_glass_requested"] = bool(request_break_glass)
|
||||||
|
payload["break_glass_authorized"] = break_glass_authorized
|
||||||
|
# Even an authorized gate performs no restart here: execution is a later
|
||||||
|
# child. The gate proves the apply path *would* be permitted.
|
||||||
|
payload["reasons"] = list(payload.get("reasons") or []) + list(
|
||||||
|
gate_payload.get("reasons") or []
|
||||||
|
)
|
||||||
|
if not gate.allow and gate.incident is not None:
|
||||||
|
payload["incident"] = gate.incident
|
||||||
return payload
|
return payload
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,907 @@
|
|||||||
|
"""Tests for the pre-restart drain proof and hard gate (#661).
|
||||||
|
|
||||||
|
Covers the acceptance criteria:
|
||||||
|
|
||||||
|
1. Restart apply without a proof fails closed.
|
||||||
|
2. A successful drain produces a verifiable proof.
|
||||||
|
3. An open unsafe mutation makes the proof fail (multi-session fixture).
|
||||||
|
4. Pass / fail / expired verification paths.
|
||||||
|
|
||||||
|
Plus the security posture: forged/tampered proofs are rejected, break-glass is
|
||||||
|
the only bypass and is never silent, a stale blast-radius fingerprint rejects a
|
||||||
|
proof, and no per-process secret ever leaks into a serialized artifact.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
import unittest
|
||||||
|
from datetime import datetime, timedelta, timezone
|
||||||
|
|
||||||
|
import drain_proof as dp
|
||||||
|
import restart_coordinator as rc
|
||||||
|
|
||||||
|
|
||||||
|
NOW = datetime(2026, 7, 24, 6, 0, 0, tzinfo=timezone.utc)
|
||||||
|
SECRET = b"unit-test-drain-proof-secret-0123456789abcdef"
|
||||||
|
|
||||||
|
|
||||||
|
def _live_pid() -> int:
|
||||||
|
return os.getpid()
|
||||||
|
|
||||||
|
|
||||||
|
def _clean_drain_state() -> dict:
|
||||||
|
"""Every drain action succeeded, no sessions outstanding."""
|
||||||
|
|
||||||
|
return {
|
||||||
|
"assignments_stopped": True,
|
||||||
|
"checkpoints_complete": True,
|
||||||
|
"handoffs_verified": True,
|
||||||
|
"leases_handled": True,
|
||||||
|
"acks": {}, # no other live sessions to acknowledge
|
||||||
|
"ack_timeout_policy_applied": False,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _safe_report() -> dict:
|
||||||
|
"""Impact report with no other live work: a restart here is safe."""
|
||||||
|
|
||||||
|
report = rc.evaluate_restart_impact(
|
||||||
|
{"sessions": [], "leases": [], "inventory_complete": True},
|
||||||
|
now=NOW,
|
||||||
|
requesting_session_id="prgs-controller-1-req",
|
||||||
|
)
|
||||||
|
return report.as_dict()
|
||||||
|
|
||||||
|
|
||||||
|
def _unsafe_mutation_report() -> dict:
|
||||||
|
"""Multi-session report: a second session holds a live author mutation."""
|
||||||
|
|
||||||
|
sessions = [
|
||||||
|
{
|
||||||
|
"session_id": "prgs-controller-1-req",
|
||||||
|
"role": "controller",
|
||||||
|
"profile": "prgs-controller",
|
||||||
|
"pid": _live_pid(),
|
||||||
|
"status": "active",
|
||||||
|
"last_heartbeat_at": NOW.isoformat(),
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"session_id": "prgs-author-99",
|
||||||
|
"role": "author",
|
||||||
|
"profile": "prgs-author",
|
||||||
|
"pid": _live_pid(),
|
||||||
|
"status": "active",
|
||||||
|
"last_heartbeat_at": NOW.isoformat(),
|
||||||
|
},
|
||||||
|
]
|
||||||
|
leases = [
|
||||||
|
{
|
||||||
|
"lease_id": "lease-mut",
|
||||||
|
"session_id": "prgs-author-99",
|
||||||
|
"role": "author",
|
||||||
|
"phase": "implementing",
|
||||||
|
"work_kind": "issue",
|
||||||
|
"work_number": 661,
|
||||||
|
"worktree_path": "branches/issue-661",
|
||||||
|
"freshness": {"freshness": "active"},
|
||||||
|
}
|
||||||
|
]
|
||||||
|
report = rc.evaluate_restart_impact(
|
||||||
|
{"sessions": sessions, "leases": leases, "inventory_complete": True},
|
||||||
|
now=NOW,
|
||||||
|
requesting_session_id="prgs-controller-1-req",
|
||||||
|
)
|
||||||
|
return report.as_dict()
|
||||||
|
|
||||||
|
|
||||||
|
class BuildDrainProofTests(unittest.TestCase):
|
||||||
|
def test_clean_drain_produces_verifiable_clean_proof(self):
|
||||||
|
"""AC#2: a successful drain produces a verifiable proof."""
|
||||||
|
|
||||||
|
proof = dp.build_drain_proof(
|
||||||
|
impact_report=_safe_report(),
|
||||||
|
drain_state=_clean_drain_state(),
|
||||||
|
requesting_session_id="prgs-controller-1-req",
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
)
|
||||||
|
self.assertTrue(proof.clean)
|
||||||
|
self.assertEqual(proof.failed_checks, [])
|
||||||
|
self.assertEqual(
|
||||||
|
{c.name for c in proof.checks}, set(dp.REQUIRED_CHECKS)
|
||||||
|
)
|
||||||
|
result = dp.verify_drain_proof(
|
||||||
|
proof.as_dict(), now=NOW, secret=SECRET
|
||||||
|
)
|
||||||
|
self.assertTrue(result.valid, result.reasons)
|
||||||
|
self.assertFalse(result.expired)
|
||||||
|
self.assertFalse(result.tampered)
|
||||||
|
|
||||||
|
def test_open_mutation_makes_proof_unclean(self):
|
||||||
|
"""AC#3: an unsafe mutation still in flight fails the proof."""
|
||||||
|
|
||||||
|
proof = dp.build_drain_proof(
|
||||||
|
impact_report=_unsafe_mutation_report(),
|
||||||
|
drain_state=_clean_drain_state(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
)
|
||||||
|
self.assertFalse(proof.clean)
|
||||||
|
self.assertIn(dp.CHECK_NO_INFLIGHT_MUTATIONS, proof.failed_checks)
|
||||||
|
# Leases-handled also fails: the report still shows a disruptive lease.
|
||||||
|
self.assertIn(dp.CHECK_LEASES_HANDLED, proof.failed_checks)
|
||||||
|
result = dp.verify_drain_proof(proof.as_dict(), now=NOW, secret=SECRET)
|
||||||
|
self.assertFalse(result.valid)
|
||||||
|
|
||||||
|
def test_incomplete_inventory_fails_no_mutations_check(self):
|
||||||
|
proof = dp.build_drain_proof(
|
||||||
|
impact_report={"inventory_complete": False},
|
||||||
|
drain_state=_clean_drain_state(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
)
|
||||||
|
self.assertFalse(proof.clean)
|
||||||
|
self.assertIn(dp.CHECK_NO_INFLIGHT_MUTATIONS, proof.failed_checks)
|
||||||
|
|
||||||
|
def test_missing_checkpoint_flag_fails_closed(self):
|
||||||
|
state = _clean_drain_state()
|
||||||
|
del state["checkpoints_complete"]
|
||||||
|
proof = dp.build_drain_proof(
|
||||||
|
impact_report=_safe_report(), drain_state=state, now=NOW, secret=SECRET
|
||||||
|
)
|
||||||
|
self.assertFalse(proof.clean)
|
||||||
|
self.assertIn(dp.CHECK_CHECKPOINTS_COMPLETE, proof.failed_checks)
|
||||||
|
|
||||||
|
def test_non_true_flags_fail_closed(self):
|
||||||
|
"""A truthy-but-not-True value (e.g. the string 'yes') must not pass."""
|
||||||
|
|
||||||
|
state = _clean_drain_state()
|
||||||
|
state["assignments_stopped"] = "yes"
|
||||||
|
proof = dp.build_drain_proof(
|
||||||
|
impact_report=_safe_report(), drain_state=state, now=NOW, secret=SECRET
|
||||||
|
)
|
||||||
|
self.assertIn(dp.CHECK_ASSIGNMENTS_STOPPED, proof.failed_checks)
|
||||||
|
|
||||||
|
def test_ack_timeout_policy_satisfies_ack_check(self):
|
||||||
|
state = _clean_drain_state()
|
||||||
|
state["acks"] = {"prgs-author-99": "pending"}
|
||||||
|
state["ack_timeout_policy_applied"] = True
|
||||||
|
proof = dp.build_drain_proof(
|
||||||
|
impact_report=_safe_report(), drain_state=state, now=NOW, secret=SECRET
|
||||||
|
)
|
||||||
|
names = {c.name: c.passed for c in proof.checks}
|
||||||
|
self.assertTrue(names[dp.CHECK_ACKS_OR_TIMEOUT])
|
||||||
|
|
||||||
|
def test_outstanding_acks_without_timeout_fail(self):
|
||||||
|
state = _clean_drain_state()
|
||||||
|
state["acks"] = {"prgs-author-99": "pending"}
|
||||||
|
state["ack_timeout_policy_applied"] = False
|
||||||
|
proof = dp.build_drain_proof(
|
||||||
|
impact_report=_safe_report(), drain_state=state, now=NOW, secret=SECRET
|
||||||
|
)
|
||||||
|
self.assertIn(dp.CHECK_ACKS_OR_TIMEOUT, proof.failed_checks)
|
||||||
|
|
||||||
|
def test_all_acked_satisfies_ack_check(self):
|
||||||
|
state = _clean_drain_state()
|
||||||
|
state["acks"] = {"prgs-author-99": "acked", "prgs-author-2": "acknowledged"}
|
||||||
|
proof = dp.build_drain_proof(
|
||||||
|
impact_report=_safe_report(), drain_state=state, now=NOW, secret=SECRET
|
||||||
|
)
|
||||||
|
names = {c.name: c.passed for c in proof.checks}
|
||||||
|
self.assertTrue(names[dp.CHECK_ACKS_OR_TIMEOUT])
|
||||||
|
|
||||||
|
|
||||||
|
class VerifyDrainProofTests(unittest.TestCase):
|
||||||
|
def _clean_proof_dict(self) -> dict:
|
||||||
|
return dp.build_drain_proof(
|
||||||
|
impact_report=_safe_report(),
|
||||||
|
drain_state=_clean_drain_state(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
).as_dict()
|
||||||
|
|
||||||
|
def test_missing_proof_is_invalid(self):
|
||||||
|
result = dp.verify_drain_proof(None, now=NOW, secret=SECRET)
|
||||||
|
self.assertFalse(result.valid)
|
||||||
|
self.assertIsNone(result.proof_id)
|
||||||
|
|
||||||
|
def test_expired_proof_is_invalid(self):
|
||||||
|
"""AC#4: an expired proof fails verification."""
|
||||||
|
|
||||||
|
proof = self._clean_proof_dict()
|
||||||
|
later = NOW + timedelta(seconds=dp.DEFAULT_PROOF_TTL_SECONDS + 1)
|
||||||
|
result = dp.verify_drain_proof(proof, now=later, secret=SECRET)
|
||||||
|
self.assertFalse(result.valid)
|
||||||
|
self.assertTrue(result.expired)
|
||||||
|
|
||||||
|
def test_proof_valid_just_before_expiry(self):
|
||||||
|
proof = self._clean_proof_dict()
|
||||||
|
almost = NOW + timedelta(seconds=dp.DEFAULT_PROOF_TTL_SECONDS - 1)
|
||||||
|
result = dp.verify_drain_proof(proof, now=almost, secret=SECRET)
|
||||||
|
self.assertTrue(result.valid, result.reasons)
|
||||||
|
|
||||||
|
def test_wrong_secret_rejected(self):
|
||||||
|
"""A proof minted in a prior process (different secret) will not verify."""
|
||||||
|
|
||||||
|
proof = self._clean_proof_dict()
|
||||||
|
result = dp.verify_drain_proof(proof, now=NOW, secret=b"other-secret")
|
||||||
|
self.assertFalse(result.valid)
|
||||||
|
self.assertTrue(result.tampered)
|
||||||
|
|
||||||
|
def test_flipping_clean_flag_is_detected(self):
|
||||||
|
"""Forging clean=True on an unclean proof breaks the signature."""
|
||||||
|
|
||||||
|
unclean = dp.build_drain_proof(
|
||||||
|
impact_report=_unsafe_mutation_report(),
|
||||||
|
drain_state=_clean_drain_state(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
).as_dict()
|
||||||
|
self.assertFalse(unclean["clean"])
|
||||||
|
unclean["clean"] = True # forge
|
||||||
|
result = dp.verify_drain_proof(unclean, now=NOW, secret=SECRET)
|
||||||
|
self.assertFalse(result.valid)
|
||||||
|
self.assertTrue(result.tampered)
|
||||||
|
|
||||||
|
def test_tampering_a_check_is_detected(self):
|
||||||
|
unclean = dp.build_drain_proof(
|
||||||
|
impact_report=_unsafe_mutation_report(),
|
||||||
|
drain_state=_clean_drain_state(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
).as_dict()
|
||||||
|
for c in unclean["checks"]:
|
||||||
|
if c["name"] == dp.CHECK_NO_INFLIGHT_MUTATIONS:
|
||||||
|
c["passed"] = True # forge the failing check to pass
|
||||||
|
result = dp.verify_drain_proof(unclean, now=NOW, secret=SECRET)
|
||||||
|
self.assertFalse(result.valid)
|
||||||
|
self.assertTrue(result.tampered)
|
||||||
|
|
||||||
|
def test_missing_required_check_rejected(self):
|
||||||
|
proof = self._clean_proof_dict()
|
||||||
|
proof["checks"] = [
|
||||||
|
c for c in proof["checks"] if c["name"] != dp.CHECK_HANDOFFS_OK
|
||||||
|
]
|
||||||
|
result = dp.verify_drain_proof(proof, now=NOW, secret=SECRET)
|
||||||
|
self.assertFalse(result.valid)
|
||||||
|
|
||||||
|
def test_stale_fingerprint_rejected(self):
|
||||||
|
proof = self._clean_proof_dict()
|
||||||
|
result = dp.verify_drain_proof(
|
||||||
|
proof,
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
expected_impact_fingerprint="deadbeef",
|
||||||
|
)
|
||||||
|
self.assertFalse(result.valid)
|
||||||
|
|
||||||
|
def test_matching_fingerprint_accepted(self):
|
||||||
|
report = _safe_report()
|
||||||
|
proof = dp.build_drain_proof(
|
||||||
|
impact_report=report,
|
||||||
|
drain_state=_clean_drain_state(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
).as_dict()
|
||||||
|
fp = dp.impact_fingerprint(report)
|
||||||
|
result = dp.verify_drain_proof(
|
||||||
|
proof, now=NOW, secret=SECRET, expected_impact_fingerprint=fp
|
||||||
|
)
|
||||||
|
self.assertTrue(result.valid, result.reasons)
|
||||||
|
|
||||||
|
|
||||||
|
class GateApplyRestartTests(unittest.TestCase):
|
||||||
|
def _clean_proof_dict(self) -> dict:
|
||||||
|
return dp.build_drain_proof(
|
||||||
|
impact_report=_safe_report(),
|
||||||
|
drain_state=_clean_drain_state(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
).as_dict()
|
||||||
|
|
||||||
|
def test_apply_without_proof_denied(self):
|
||||||
|
"""AC#1: restart apply without a proof fails closed + raises incident."""
|
||||||
|
|
||||||
|
decision = dp.gate_apply_restart(proof=None, now=NOW, secret=SECRET)
|
||||||
|
self.assertFalse(decision.allow)
|
||||||
|
self.assertEqual(decision.verdict, dp.GATE_DENY)
|
||||||
|
self.assertIsNotNone(decision.incident)
|
||||||
|
self.assertEqual(
|
||||||
|
decision.incident["kind"], "restart_drain_gate_denied"
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_apply_with_valid_proof_allowed(self):
|
||||||
|
decision = dp.gate_apply_restart(
|
||||||
|
proof=self._clean_proof_dict(), now=NOW, secret=SECRET
|
||||||
|
)
|
||||||
|
self.assertTrue(decision.allow)
|
||||||
|
self.assertEqual(decision.verdict, dp.GATE_ALLOW)
|
||||||
|
self.assertIsNone(decision.incident)
|
||||||
|
|
||||||
|
def test_apply_with_expired_proof_denied_with_incident(self):
|
||||||
|
later = NOW + timedelta(seconds=dp.DEFAULT_PROOF_TTL_SECONDS + 5)
|
||||||
|
decision = dp.gate_apply_restart(
|
||||||
|
proof=self._clean_proof_dict(), now=later, secret=SECRET
|
||||||
|
)
|
||||||
|
self.assertFalse(decision.allow)
|
||||||
|
self.assertIsNotNone(decision.incident)
|
||||||
|
|
||||||
|
def test_apply_with_unclean_proof_denied(self):
|
||||||
|
"""AC#3 at the gate: an unsafe-mutation proof is denied."""
|
||||||
|
|
||||||
|
unclean = dp.build_drain_proof(
|
||||||
|
impact_report=_unsafe_mutation_report(),
|
||||||
|
drain_state=_clean_drain_state(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
).as_dict()
|
||||||
|
decision = dp.gate_apply_restart(proof=unclean, now=NOW, secret=SECRET)
|
||||||
|
self.assertFalse(decision.allow)
|
||||||
|
self.assertIsNotNone(decision.incident)
|
||||||
|
|
||||||
|
def test_break_glass_allows_without_proof_but_records_bypass(self):
|
||||||
|
decision = dp.gate_apply_restart(
|
||||||
|
proof=None, now=NOW, secret=SECRET, break_glass=True
|
||||||
|
)
|
||||||
|
self.assertTrue(decision.allow)
|
||||||
|
self.assertEqual(decision.verdict, dp.GATE_BREAK_GLASS)
|
||||||
|
self.assertTrue(decision.break_glass)
|
||||||
|
self.assertIsNone(decision.incident)
|
||||||
|
self.assertTrue(decision.audit_record["break_glass"])
|
||||||
|
|
||||||
|
def test_denied_gate_carries_stale_fingerprint_reason(self):
|
||||||
|
decision = dp.gate_apply_restart(
|
||||||
|
proof=self._clean_proof_dict(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
expected_impact_fingerprint="not-the-fingerprint",
|
||||||
|
)
|
||||||
|
self.assertFalse(decision.allow)
|
||||||
|
|
||||||
|
|
||||||
|
class SecretHygieneTests(unittest.TestCase):
|
||||||
|
def test_secret_never_serialized(self):
|
||||||
|
proof = dp.build_drain_proof(
|
||||||
|
impact_report=_safe_report(),
|
||||||
|
drain_state=_clean_drain_state(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
)
|
||||||
|
blob = dp._canonical(proof.as_dict())
|
||||||
|
self.assertNotIn(SECRET.decode(), blob)
|
||||||
|
# The signature is a hex digest, not the raw secret.
|
||||||
|
self.assertNotIn(SECRET.hex(), blob)
|
||||||
|
|
||||||
|
def test_incident_descriptor_has_no_secret(self):
|
||||||
|
decision = dp.gate_apply_restart(proof=None, now=NOW, secret=SECRET)
|
||||||
|
blob = dp._canonical(decision.incident)
|
||||||
|
self.assertNotIn(SECRET.decode(), blob)
|
||||||
|
|
||||||
|
|
||||||
|
def _drained_report_with_live_sessions(count: int) -> dict:
|
||||||
|
"""Report with ``count`` other live sessions but nothing in flight.
|
||||||
|
|
||||||
|
Every other checklist item passes against this report, so a failure
|
||||||
|
isolates the acknowledgement check rather than tripping on mutations.
|
||||||
|
"""
|
||||||
|
|
||||||
|
sessions = [
|
||||||
|
{
|
||||||
|
"session_id": "prgs-controller-1-req",
|
||||||
|
"role": "controller",
|
||||||
|
"profile": "prgs-controller",
|
||||||
|
"pid": _live_pid(),
|
||||||
|
"status": "active",
|
||||||
|
"last_heartbeat_at": NOW.isoformat(),
|
||||||
|
}
|
||||||
|
]
|
||||||
|
for index in range(count):
|
||||||
|
sessions.append(
|
||||||
|
{
|
||||||
|
"session_id": f"prgs-author-{index}",
|
||||||
|
"role": "author",
|
||||||
|
"profile": "prgs-author",
|
||||||
|
"pid": _live_pid(),
|
||||||
|
"status": "active",
|
||||||
|
"last_heartbeat_at": NOW.isoformat(),
|
||||||
|
}
|
||||||
|
)
|
||||||
|
report = rc.evaluate_restart_impact(
|
||||||
|
{"sessions": sessions, "leases": [], "inventory_complete": True},
|
||||||
|
now=NOW,
|
||||||
|
requesting_session_id="prgs-controller-1-req",
|
||||||
|
)
|
||||||
|
return report.as_dict()
|
||||||
|
|
||||||
|
|
||||||
|
class AcknowledgementFailClosedTests(unittest.TestCase):
|
||||||
|
"""Acknowledgement evidence must fail closed unless explicitly verified.
|
||||||
|
|
||||||
|
Regression cover for the reviewed fail-open on PR #882: an absent ``acks``
|
||||||
|
key collapsed to ``{}`` and was read as "no other live sessions required to
|
||||||
|
acknowledge", so a proof minted clean and the restart gate allowed while the
|
||||||
|
impact report still showed other live sessions.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def _state(self, **overrides) -> dict:
|
||||||
|
state = _clean_drain_state()
|
||||||
|
state.pop("acks", None)
|
||||||
|
state["ack_timeout_policy_applied"] = False
|
||||||
|
state.update(overrides)
|
||||||
|
return state
|
||||||
|
|
||||||
|
def _acks_check(self, proof) -> dp.DrainCheck:
|
||||||
|
return next(c for c in proof.checks if c.name == dp.CHECK_ACKS_OR_TIMEOUT)
|
||||||
|
|
||||||
|
def _build(self, report: dict, state: dict):
|
||||||
|
return dp.build_drain_proof(
|
||||||
|
impact_report=report, drain_state=state, now=NOW, secret=SECRET
|
||||||
|
)
|
||||||
|
|
||||||
|
def assertAcksFailClosed(self, report: dict, state: dict) -> None:
|
||||||
|
proof = self._build(report, state)
|
||||||
|
self.assertFalse(self._acks_check(proof).passed)
|
||||||
|
self.assertIn(dp.CHECK_ACKS_OR_TIMEOUT, proof.failed_checks)
|
||||||
|
self.assertFalse(proof.clean)
|
||||||
|
|
||||||
|
# --- missing / null / empty / malformed ------------------------------
|
||||||
|
|
||||||
|
def test_missing_acks_key_with_live_sessions_fails_closed(self):
|
||||||
|
"""The exact reviewed defect: absent key, three other live sessions."""
|
||||||
|
report = _drained_report_with_live_sessions(3)
|
||||||
|
self.assertEqual(report["counts"]["sessions_live_other"], 3)
|
||||||
|
state = self._state()
|
||||||
|
self.assertNotIn("acks", state)
|
||||||
|
proof = self._build(report, state)
|
||||||
|
check = self._acks_check(proof)
|
||||||
|
self.assertFalse(check.passed)
|
||||||
|
self.assertNotIn("no other live sessions", check.detail)
|
||||||
|
self.assertIn("fail closed", check.detail)
|
||||||
|
self.assertFalse(proof.clean)
|
||||||
|
self.assertEqual(proof.failed_checks, [dp.CHECK_ACKS_OR_TIMEOUT])
|
||||||
|
|
||||||
|
def test_none_acks_with_live_sessions_fails_closed(self):
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
_drained_report_with_live_sessions(2), self._state(acks=None)
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_empty_acks_with_live_sessions_fails_closed(self):
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
_drained_report_with_live_sessions(1), self._state(acks={})
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_malformed_acks_fail_closed(self):
|
||||||
|
for malformed in ([], "ack", 7, ("ack",), True):
|
||||||
|
with self.subTest(malformed=malformed):
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
_drained_report_with_live_sessions(1),
|
||||||
|
self._state(acks=malformed),
|
||||||
|
)
|
||||||
|
|
||||||
|
# --- stale / unproven values -----------------------------------------
|
||||||
|
|
||||||
|
def test_stale_or_unproven_ack_values_fail_closed(self):
|
||||||
|
for value in ("pending", "stale", "unknown", "", None, True, 1, NOW):
|
||||||
|
with self.subTest(value=value):
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
_drained_report_with_live_sessions(1),
|
||||||
|
self._state(acks={"prgs-author-0": value}),
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_partial_coverage_fails_closed(self):
|
||||||
|
"""Fewer acknowledgements than the report's live-session count."""
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
_drained_report_with_live_sessions(3),
|
||||||
|
self._state(acks={"prgs-author-0": "ack"}),
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_one_unacked_entry_among_many_fails_closed(self):
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
_drained_report_with_live_sessions(2),
|
||||||
|
self._state(acks={"prgs-author-0": "ack", "prgs-author-1": "pending"}),
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_unproven_live_session_count_fails_closed(self):
|
||||||
|
"""A missing or malformed count cannot prove nobody had to acknowledge."""
|
||||||
|
malformed_counts = (
|
||||||
|
None,
|
||||||
|
{},
|
||||||
|
{"sessions_live_other": None},
|
||||||
|
{"sessions_live_other": "3"},
|
||||||
|
{"sessions_live_other": -1},
|
||||||
|
{"sessions_live_other": True},
|
||||||
|
)
|
||||||
|
for counts in malformed_counts:
|
||||||
|
with self.subTest(counts=counts):
|
||||||
|
report = _drained_report_with_live_sessions(0)
|
||||||
|
if counts is None:
|
||||||
|
report.pop("counts", None)
|
||||||
|
else:
|
||||||
|
report["counts"] = counts
|
||||||
|
self.assertAcksFailClosed(report, self._state())
|
||||||
|
|
||||||
|
# --- valid evidence still passes -------------------------------------
|
||||||
|
|
||||||
|
def test_complete_valid_acks_pass(self):
|
||||||
|
report = _drained_report_with_live_sessions(2)
|
||||||
|
state = self._state(
|
||||||
|
acks={"prgs-author-0": "ack", "prgs-author-1": "acknowledged"}
|
||||||
|
)
|
||||||
|
proof = self._build(report, state)
|
||||||
|
self.assertTrue(self._acks_check(proof).passed)
|
||||||
|
self.assertTrue(proof.clean)
|
||||||
|
self.assertEqual(proof.failed_checks, [])
|
||||||
|
|
||||||
|
def test_no_other_live_sessions_still_passes(self):
|
||||||
|
"""Intended behavior retained: zero live sessions needs no acks."""
|
||||||
|
report = _drained_report_with_live_sessions(0)
|
||||||
|
self.assertEqual(report["counts"]["sessions_live_other"], 0)
|
||||||
|
proof = self._build(report, self._state())
|
||||||
|
check = self._acks_check(proof)
|
||||||
|
self.assertTrue(check.passed)
|
||||||
|
self.assertIn("sessions_live_other=0", check.detail)
|
||||||
|
self.assertTrue(proof.clean)
|
||||||
|
|
||||||
|
# --- timeout policy cannot become a second fail-open ------------------
|
||||||
|
|
||||||
|
def test_unproven_timeout_policy_cannot_open_the_gate(self):
|
||||||
|
for value in (None, "true", "yes", 1, "True", [], {}):
|
||||||
|
with self.subTest(value=value):
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
_drained_report_with_live_sessions(2),
|
||||||
|
self._state(ack_timeout_policy_applied=value),
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_explicit_timeout_policy_permits(self):
|
||||||
|
proof = self._build(
|
||||||
|
_drained_report_with_live_sessions(2),
|
||||||
|
self._state(ack_timeout_policy_applied=True),
|
||||||
|
)
|
||||||
|
check = self._acks_check(proof)
|
||||||
|
self.assertTrue(check.passed)
|
||||||
|
self.assertIn("timeout policy", check.detail)
|
||||||
|
self.assertTrue(proof.clean)
|
||||||
|
|
||||||
|
# --- the gate itself must deny ---------------------------------------
|
||||||
|
|
||||||
|
def test_failed_ack_check_denies_the_restart_gate(self):
|
||||||
|
report = _drained_report_with_live_sessions(3)
|
||||||
|
proof = self._build(report, self._state())
|
||||||
|
self.assertFalse(proof.clean)
|
||||||
|
decision = dp.gate_apply_restart(
|
||||||
|
proof=proof.as_dict(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
expected_impact_fingerprint=dp.impact_fingerprint(report),
|
||||||
|
)
|
||||||
|
self.assertFalse(decision.allow)
|
||||||
|
self.assertEqual(decision.verdict, dp.GATE_DENY)
|
||||||
|
self.assertIsNotNone(decision.incident)
|
||||||
|
|
||||||
|
def test_unclean_ack_proof_fails_verification(self):
|
||||||
|
report = _drained_report_with_live_sessions(3)
|
||||||
|
proof = self._build(report, self._state())
|
||||||
|
result = dp.verify_drain_proof(
|
||||||
|
proof.as_dict(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
expected_impact_fingerprint=dp.impact_fingerprint(report),
|
||||||
|
)
|
||||||
|
self.assertFalse(result.valid)
|
||||||
|
self.assertFalse(result.clean)
|
||||||
|
|
||||||
|
|
||||||
|
def _identity_report(*, requester: str, others: tuple[str, ...]) -> dict:
|
||||||
|
"""Report with explicitly named requester and other live sessions.
|
||||||
|
|
||||||
|
Unlike :func:`_drained_report_with_live_sessions`, the session ids are
|
||||||
|
chosen by the caller so a test can supply acknowledgements for the *wrong*
|
||||||
|
identities while keeping the count correct.
|
||||||
|
"""
|
||||||
|
|
||||||
|
sessions = [
|
||||||
|
{
|
||||||
|
"session_id": requester,
|
||||||
|
"role": "controller",
|
||||||
|
"profile": "prgs-controller",
|
||||||
|
"pid": _live_pid(),
|
||||||
|
"status": "active",
|
||||||
|
"last_heartbeat_at": NOW.isoformat(),
|
||||||
|
}
|
||||||
|
]
|
||||||
|
for session_id in others:
|
||||||
|
sessions.append(
|
||||||
|
{
|
||||||
|
"session_id": session_id,
|
||||||
|
"role": "author",
|
||||||
|
"profile": "prgs-author",
|
||||||
|
"pid": _live_pid(),
|
||||||
|
"status": "active",
|
||||||
|
"last_heartbeat_at": NOW.isoformat(),
|
||||||
|
}
|
||||||
|
)
|
||||||
|
report = rc.evaluate_restart_impact(
|
||||||
|
{"sessions": sessions, "leases": [], "inventory_complete": True},
|
||||||
|
now=NOW,
|
||||||
|
requesting_session_id=requester,
|
||||||
|
)
|
||||||
|
return report.as_dict()
|
||||||
|
|
||||||
|
|
||||||
|
class AcknowledgementIdentityBindingTests(unittest.TestCase):
|
||||||
|
"""Acknowledgement coverage must be bound to session identity, not counted.
|
||||||
|
|
||||||
|
Regression cover for the second reviewed fail-open on PR #882 (review 582,
|
||||||
|
blocker B1): coverage compared ``acked_count`` against
|
||||||
|
``counts.sessions_live_other``, so acknowledgements supplied for the
|
||||||
|
requesting session and for ids that do not exist satisfied the obligations
|
||||||
|
of the live sessions that never answered. The required identities are
|
||||||
|
carried by the report itself — ``ack_state`` keys and ``affected_sessions``
|
||||||
|
filtered on ``live and not is_requester`` — and only an acknowledgement
|
||||||
|
keyed by one of those ids may count for it.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def _state(self, **overrides) -> dict:
|
||||||
|
state = _clean_drain_state()
|
||||||
|
state.pop("acks", None)
|
||||||
|
state["ack_timeout_policy_applied"] = False
|
||||||
|
state.update(overrides)
|
||||||
|
return state
|
||||||
|
|
||||||
|
def _acks_check(self, proof) -> dp.DrainCheck:
|
||||||
|
return next(c for c in proof.checks if c.name == dp.CHECK_ACKS_OR_TIMEOUT)
|
||||||
|
|
||||||
|
def _build(self, report: dict, state: dict):
|
||||||
|
return dp.build_drain_proof(
|
||||||
|
impact_report=report, drain_state=state, now=NOW, secret=SECRET
|
||||||
|
)
|
||||||
|
|
||||||
|
def assertAcksFailClosed(self, report: dict, state: dict) -> dp.DrainCheck:
|
||||||
|
"""Failure must propagate through the check, the proof, and the gate."""
|
||||||
|
|
||||||
|
proof = self._build(report, state)
|
||||||
|
check = self._acks_check(proof)
|
||||||
|
self.assertFalse(check.passed)
|
||||||
|
self.assertFalse(proof.clean)
|
||||||
|
self.assertIn(dp.CHECK_ACKS_OR_TIMEOUT, proof.failed_checks)
|
||||||
|
decision = dp.gate_apply_restart(
|
||||||
|
proof=proof.as_dict(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
expected_impact_fingerprint=dp.impact_fingerprint(report),
|
||||||
|
)
|
||||||
|
self.assertEqual(decision.verdict, dp.GATE_DENY)
|
||||||
|
self.assertFalse(decision.allow)
|
||||||
|
return check
|
||||||
|
|
||||||
|
# --- the reviewer's exact reproduction --------------------------------
|
||||||
|
|
||||||
|
def test_requester_plus_unknown_id_cannot_satisfy_two_live_sessions(self):
|
||||||
|
"""Review 582 B1 verbatim: requester + a nonexistent session.
|
||||||
|
|
||||||
|
``sessions_live_other=2`` with ``ack_state`` naming ``other-0`` and
|
||||||
|
``other-1``; the drain state supplies an acknowledgement from the
|
||||||
|
requesting session itself and from a session that does not exist. The
|
||||||
|
count matches, the identities do not.
|
||||||
|
"""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0", "other-1"))
|
||||||
|
self.assertEqual(report["counts"]["sessions_live_other"], 2)
|
||||||
|
self.assertEqual(
|
||||||
|
report["ack_state"], {"other-0": "pending", "other-1": "pending"}
|
||||||
|
)
|
||||||
|
state = self._state(acks={"req": "ack", "totally-bogus-session": "ack"})
|
||||||
|
check = self.assertAcksFailClosed(report, state)
|
||||||
|
self.assertIn("other-0", check.detail)
|
||||||
|
self.assertIn("other-1", check.detail)
|
||||||
|
self.assertIn("fail closed", check.detail)
|
||||||
|
|
||||||
|
# --- wrong / unknown / requester identities ---------------------------
|
||||||
|
|
||||||
|
def test_sufficient_count_of_wrong_ids_fails_closed(self):
|
||||||
|
"""Right cardinality, wrong identities: two acks, neither required."""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0", "other-1"))
|
||||||
|
state = self._state(acks={"ghost-a": "ack", "ghost-b": "ack"})
|
||||||
|
check = self.assertAcksFailClosed(report, state)
|
||||||
|
self.assertIn("do not count", check.detail)
|
||||||
|
|
||||||
|
def test_more_acks_than_required_still_fails_on_wrong_ids(self):
|
||||||
|
"""Coverage cannot be bought with volume: five acks, none required."""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0", "other-1"))
|
||||||
|
state = self._state(acks={f"ghost-{i}": "acknowledged" for i in range(5)})
|
||||||
|
self.assertAcksFailClosed(report, state)
|
||||||
|
|
||||||
|
def test_partial_identity_match_fails_closed(self):
|
||||||
|
"""One required id acknowledged, the rest padded with unknown ids."""
|
||||||
|
|
||||||
|
report = _identity_report(
|
||||||
|
requester="req", others=("other-0", "other-1", "other-2")
|
||||||
|
)
|
||||||
|
state = self._state(
|
||||||
|
acks={"other-0": "ack", "ghost-1": "ack", "ghost-2": "ack"}
|
||||||
|
)
|
||||||
|
check = self.assertAcksFailClosed(report, state)
|
||||||
|
self.assertIn("other-1", check.detail)
|
||||||
|
self.assertIn("other-2", check.detail)
|
||||||
|
|
||||||
|
def test_requester_ack_never_satisfies_another_sessions_obligation(self):
|
||||||
|
"""The requester is excluded from the required set and stays excluded."""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
requester_rows = [s for s in report["affected_sessions"] if s["is_requester"]]
|
||||||
|
self.assertEqual([s["session_id"] for s in requester_rows], ["req"])
|
||||||
|
self.assertNotIn("req", report["ack_state"])
|
||||||
|
check = self.assertAcksFailClosed(report, self._state(acks={"req": "ack"}))
|
||||||
|
self.assertIn("other-0", check.detail)
|
||||||
|
|
||||||
|
def test_fabricated_ids_do_not_count_toward_coverage(self):
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
for bogus in ("", " ", "other-0 extra", "OTHER-0", "other-01", "0"):
|
||||||
|
with self.subTest(bogus=bogus):
|
||||||
|
self.assertAcksFailClosed(report, self._state(acks={bogus: "ack"}))
|
||||||
|
|
||||||
|
# --- per-session state must be explicitly valid ------------------------
|
||||||
|
|
||||||
|
def test_unproven_per_session_states_fail_closed(self):
|
||||||
|
"""A required id present but not explicitly acknowledged fails closed."""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0", "other-1"))
|
||||||
|
for value in ("pending", "stale", "unknown", "", None, True, 1, NOW):
|
||||||
|
with self.subTest(value=value):
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
report,
|
||||||
|
self._state(acks={"other-0": "ack", "other-1": value}),
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_report_ack_state_placeholder_is_never_read_as_an_ack(self):
|
||||||
|
"""``ack_state`` values are the report's own placeholders, not evidence."""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
report["ack_state"] = {"other-0": "ack"}
|
||||||
|
self.assertAcksFailClosed(report, self._state())
|
||||||
|
|
||||||
|
# --- missing / malformed / contradictory identity evidence -------------
|
||||||
|
|
||||||
|
def test_missing_identity_evidence_fails_closed(self):
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
report.pop("ack_state", None)
|
||||||
|
report.pop("affected_sessions", None)
|
||||||
|
check = self.assertAcksFailClosed(report, self._state(acks={"other-0": "ack"}))
|
||||||
|
self.assertIn("no session-identity evidence", check.detail)
|
||||||
|
|
||||||
|
def test_malformed_ack_state_fails_closed(self):
|
||||||
|
for malformed in ([], "other-0", 7, None, ("other-0",)):
|
||||||
|
with self.subTest(malformed=malformed):
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
report["ack_state"] = malformed
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
report, self._state(acks={"other-0": "ack"})
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_non_string_ack_state_key_fails_closed(self):
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
report["ack_state"] = {7: "pending"}
|
||||||
|
self.assertAcksFailClosed(report, self._state(acks={"other-0": "ack"}))
|
||||||
|
|
||||||
|
def test_malformed_affected_sessions_fails_closed(self):
|
||||||
|
for malformed in ("sessions", 7, {"session_id": "other-0"}, [None], [7]):
|
||||||
|
with self.subTest(malformed=malformed):
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
report.pop("ack_state", None)
|
||||||
|
report["affected_sessions"] = malformed
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
report, self._state(acks={"other-0": "ack"})
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_affected_sessions_without_explicit_booleans_fails_closed(self):
|
||||||
|
"""``live``/``is_requester`` must be real booleans, never inferred."""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
report.pop("ack_state", None)
|
||||||
|
for row in report["affected_sessions"]:
|
||||||
|
if row["session_id"] == "other-0":
|
||||||
|
row["is_requester"] = "false"
|
||||||
|
self.assertAcksFailClosed(report, self._state(acks={"other-0": "ack"}))
|
||||||
|
|
||||||
|
def test_affected_sessions_missing_live_flag_fails_closed(self):
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
report.pop("ack_state", None)
|
||||||
|
for row in report["affected_sessions"]:
|
||||||
|
row.pop("live", None)
|
||||||
|
self.assertAcksFailClosed(report, self._state(acks={"other-0": "ack"}))
|
||||||
|
|
||||||
|
def test_contradictory_ack_state_and_affected_sessions_fails_closed(self):
|
||||||
|
"""Both views present and disagreeing is unresolvable, not a tie-break."""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0", "other-1"))
|
||||||
|
report["ack_state"] = {"other-0": "pending", "other-9": "pending"}
|
||||||
|
check = self.assertAcksFailClosed(
|
||||||
|
report, self._state(acks={"other-0": "ack", "other-9": "ack"})
|
||||||
|
)
|
||||||
|
self.assertIn("contradicts itself", check.detail)
|
||||||
|
|
||||||
|
def test_identity_count_mismatch_fails_closed(self):
|
||||||
|
"""Identity evidence that cannot be reconciled with the count denies."""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0", "other-1"))
|
||||||
|
report["counts"] = dict(report["counts"], sessions_live_other=1)
|
||||||
|
check = self.assertAcksFailClosed(
|
||||||
|
report, self._state(acks={"other-0": "ack", "other-1": "ack"})
|
||||||
|
)
|
||||||
|
self.assertIn("cannot be reconciled", check.detail)
|
||||||
|
|
||||||
|
def test_broken_identity_evidence_outranks_timeout_policy(self):
|
||||||
|
"""The sanctioned timeout path cannot paper over an unreadable report."""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
report["ack_state"] = "not-a-mapping"
|
||||||
|
self.assertAcksFailClosed(report, self._state(ack_timeout_policy_applied=True))
|
||||||
|
|
||||||
|
# --- legitimate success is preserved -----------------------------------
|
||||||
|
|
||||||
|
def test_every_required_session_acknowledged_passes(self):
|
||||||
|
report = _identity_report(
|
||||||
|
requester="req", others=("other-0", "other-1", "other-2")
|
||||||
|
)
|
||||||
|
state = self._state(
|
||||||
|
acks={
|
||||||
|
"other-0": "ack",
|
||||||
|
"other-1": "acked",
|
||||||
|
"other-2": "acknowledged",
|
||||||
|
}
|
||||||
|
)
|
||||||
|
proof = self._build(report, state)
|
||||||
|
check = self._acks_check(proof)
|
||||||
|
self.assertTrue(check.passed)
|
||||||
|
self.assertTrue(proof.clean)
|
||||||
|
self.assertEqual(proof.failed_checks, [])
|
||||||
|
self.assertIn("acknowledged by identity", check.detail)
|
||||||
|
decision = dp.gate_apply_restart(
|
||||||
|
proof=proof.as_dict(),
|
||||||
|
now=NOW,
|
||||||
|
secret=SECRET,
|
||||||
|
expected_impact_fingerprint=dp.impact_fingerprint(report),
|
||||||
|
)
|
||||||
|
self.assertEqual(decision.verdict, dp.GATE_ALLOW)
|
||||||
|
self.assertTrue(decision.allow)
|
||||||
|
|
||||||
|
def test_required_session_ack_tolerates_surrounding_whitespace(self):
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
proof = self._build(report, self._state(acks={" other-0 ": " ACK "}))
|
||||||
|
self.assertTrue(self._acks_check(proof).passed)
|
||||||
|
self.assertTrue(proof.clean)
|
||||||
|
|
||||||
|
def test_no_other_live_sessions_still_passes_with_identity_evidence(self):
|
||||||
|
report = _identity_report(requester="req", others=())
|
||||||
|
self.assertEqual(report["counts"]["sessions_live_other"], 0)
|
||||||
|
self.assertEqual(report["ack_state"], {})
|
||||||
|
proof = self._build(report, self._state())
|
||||||
|
check = self._acks_check(proof)
|
||||||
|
self.assertTrue(check.passed)
|
||||||
|
self.assertIn("sessions_live_other=0", check.detail)
|
||||||
|
self.assertTrue(proof.clean)
|
||||||
|
|
||||||
|
def test_explicit_timeout_policy_retains_intended_behavior(self):
|
||||||
|
"""Valid, correctly typed timeout evidence still permits the check."""
|
||||||
|
|
||||||
|
report = _identity_report(requester="req", others=("other-0", "other-1"))
|
||||||
|
proof = self._build(report, self._state(ack_timeout_policy_applied=True))
|
||||||
|
check = self._acks_check(proof)
|
||||||
|
self.assertTrue(check.passed)
|
||||||
|
self.assertIn("timeout policy", check.detail)
|
||||||
|
self.assertTrue(proof.clean)
|
||||||
|
|
||||||
|
def test_timeout_policy_still_strictly_typed_under_identity_binding(self):
|
||||||
|
report = _identity_report(requester="req", others=("other-0",))
|
||||||
|
for value in (None, "true", "True", 1, [], {}):
|
||||||
|
with self.subTest(value=value):
|
||||||
|
self.assertAcksFailClosed(
|
||||||
|
report, self._state(ack_timeout_policy_applied=value)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
@@ -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,453 @@
|
|||||||
|
"""#897: stale-runtime / runtime-mode refusals must not look like permission denials.
|
||||||
|
|
||||||
|
Acceptance criteria (issue #897):
|
||||||
|
|
||||||
|
* Stale-runtime and runtime-mode refusals are typed distinctly from
|
||||||
|
profile-permission refusals (distinct ``blocker_kind``).
|
||||||
|
* A refusal caused by staleness or runtime mode never emits a
|
||||||
|
``permission_report`` and never names a permission the active profile holds.
|
||||||
|
* ``_permission_block_report`` verifies the active profile actually lacks the
|
||||||
|
operation before reporting it missing.
|
||||||
|
* A stale-runtime refusal reports reconnect-only recovery and never recommends
|
||||||
|
``gitea_activate_profile`` or an MCP session switch.
|
||||||
|
* The blocker payload states the observed heads (parity fields).
|
||||||
|
* Matrix across author / reviewer / merger / reconciler profiles.
|
||||||
|
* Regression: ``gitea_create_issue`` on a stale daemon under ``prgs-author``
|
||||||
|
never returns ``missing_permission: gitea.issue.create``.
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
import unittest
|
||||||
|
from unittest.mock import patch
|
||||||
|
|
||||||
|
sys.path.insert(0, str(__import__("pathlib").Path(__file__).resolve().parent.parent))
|
||||||
|
|
||||||
|
import gitea_config # noqa: E402
|
||||||
|
import gitea_mcp_server as mcp_server # noqa: E402
|
||||||
|
|
||||||
|
SHA_START = "7af40fb5ff7debd5e9165fe97d9c7c279358e175"
|
||||||
|
SHA_LIVE = "2f4dec832327513118f2fe92b74da25d124a01cb"
|
||||||
|
|
||||||
|
ROLE_MATRIX = (
|
||||||
|
(
|
||||||
|
"prgs-author",
|
||||||
|
"author",
|
||||||
|
"gitea.issue.create",
|
||||||
|
[
|
||||||
|
"gitea.read",
|
||||||
|
"gitea.issue.create",
|
||||||
|
"gitea.issue.comment",
|
||||||
|
"gitea.issue.close",
|
||||||
|
"gitea.branch.create",
|
||||||
|
"gitea.branch.push",
|
||||||
|
"gitea.pr.create",
|
||||||
|
"gitea.pr.comment",
|
||||||
|
"gitea.repo.commit",
|
||||||
|
],
|
||||||
|
["gitea.pr.approve", "gitea.pr.merge", "gitea.pr.request_changes"],
|
||||||
|
"gitea.pr.merge", # forbidden op for pure-permission case
|
||||||
|
),
|
||||||
|
(
|
||||||
|
"prgs-reviewer",
|
||||||
|
"reviewer",
|
||||||
|
"gitea.pr.review",
|
||||||
|
[
|
||||||
|
"gitea.read",
|
||||||
|
"gitea.pr.review",
|
||||||
|
"gitea.pr.approve",
|
||||||
|
"gitea.pr.request_changes",
|
||||||
|
"gitea.pr.comment",
|
||||||
|
"gitea.issue.comment",
|
||||||
|
],
|
||||||
|
["gitea.branch.push", "gitea.pr.create"],
|
||||||
|
"gitea.branch.push",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
"prgs-merger",
|
||||||
|
"merger",
|
||||||
|
"gitea.pr.merge",
|
||||||
|
[
|
||||||
|
"gitea.read",
|
||||||
|
"gitea.pr.merge",
|
||||||
|
"gitea.pr.comment",
|
||||||
|
"gitea.issue.comment",
|
||||||
|
],
|
||||||
|
["gitea.pr.approve", "gitea.branch.push", "gitea.pr.create"],
|
||||||
|
"gitea.branch.push",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
"prgs-reconciler",
|
||||||
|
"reconciler",
|
||||||
|
"gitea.branch.delete",
|
||||||
|
[
|
||||||
|
"gitea.read",
|
||||||
|
"gitea.branch.delete",
|
||||||
|
"gitea.pr.comment",
|
||||||
|
"gitea.issue.comment",
|
||||||
|
"gitea.pr.close",
|
||||||
|
"gitea.issue.close",
|
||||||
|
],
|
||||||
|
["gitea.pr.approve", "gitea.pr.merge"],
|
||||||
|
"gitea.pr.merge",
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _profile(name: str, role: str, allowed: list[str], forbidden: list[str]) -> dict:
|
||||||
|
return {
|
||||||
|
"profile_name": name,
|
||||||
|
"role": role,
|
||||||
|
"role_kind": role,
|
||||||
|
"allowed_operations": list(allowed),
|
||||||
|
"forbidden_operations": list(forbidden),
|
||||||
|
"identity": "test-user",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _config(profiles: dict) -> dict:
|
||||||
|
return {
|
||||||
|
"version": 2,
|
||||||
|
"profiles": {
|
||||||
|
name: {
|
||||||
|
"role": p["role"],
|
||||||
|
"allowed_operations": p["allowed_operations"],
|
||||||
|
"forbidden_operations": p["forbidden_operations"],
|
||||||
|
}
|
||||||
|
for name, p in profiles.items()
|
||||||
|
},
|
||||||
|
"rules": {"allow_runtime_switching": True},
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
class Issue897Helpers(unittest.TestCase):
|
||||||
|
def test_classify_stale_reason_strings(self):
|
||||||
|
stale = (
|
||||||
|
f"live remote master is {SHA_LIVE[:12]} but the MCP server started "
|
||||||
|
f"at {SHA_START[:12]}; the daemon is stale relative to live master "
|
||||||
|
"-- restart/reconnect before mutating"
|
||||||
|
)
|
||||||
|
classified = mcp_server._classify_operation_gate_reasons([stale])
|
||||||
|
self.assertEqual(classified["stale_runtime"], [stale])
|
||||||
|
self.assertEqual(classified["permission"], [])
|
||||||
|
self.assertEqual(classified["runtime_mode"], [])
|
||||||
|
|
||||||
|
def test_classify_permission_reason(self):
|
||||||
|
reason = "profile is not allowed to gitea.pr.merge"
|
||||||
|
classified = mcp_server._classify_operation_gate_reasons([reason])
|
||||||
|
self.assertEqual(classified["permission"], [reason])
|
||||||
|
self.assertEqual(classified["stale_runtime"], [])
|
||||||
|
|
||||||
|
def test_classify_runtime_mode_reason(self):
|
||||||
|
reason = (
|
||||||
|
"runtime mode is 'dev-test' and the mutation targets the "
|
||||||
|
"production repository; dev/test runtimes must not mutate real "
|
||||||
|
"issues or PRs (ADR: stable control runtime vs dev runtime)"
|
||||||
|
)
|
||||||
|
classified = mcp_server._classify_operation_gate_reasons([reason])
|
||||||
|
self.assertEqual(classified["runtime_mode"], [reason])
|
||||||
|
self.assertEqual(classified["stale_runtime"], [])
|
||||||
|
|
||||||
|
|
||||||
|
class Issue897PermissionBlockReport(unittest.TestCase):
|
||||||
|
def test_holds_op_is_diagnostic_defect_not_missing_permission(self):
|
||||||
|
profile = _profile(
|
||||||
|
"prgs-author",
|
||||||
|
"author",
|
||||||
|
["gitea.read", "gitea.issue.create", "gitea.issue.comment"],
|
||||||
|
[],
|
||||||
|
)
|
||||||
|
with patch.object(mcp_server, "get_profile", return_value=profile), patch.object(
|
||||||
|
mcp_server.gitea_config, "load_config", return_value=_config({"prgs-author": profile})
|
||||||
|
), patch.object(
|
||||||
|
mcp_server.gitea_config, "is_runtime_switching_enabled", return_value=True
|
||||||
|
):
|
||||||
|
report = mcp_server._permission_block_report("gitea.issue.create")
|
||||||
|
self.assertTrue(report.get("diagnostic_defect"), report)
|
||||||
|
self.assertIsNone(report.get("missing_permission"), report)
|
||||||
|
action = (report.get("exact_safe_next_action") or "").lower()
|
||||||
|
# Must not *recommend* profile switching; mentioning the forbidden
|
||||||
|
# action in a "do not call" instruction is fine.
|
||||||
|
self.assertNotIn("call gitea_activate_profile with", action)
|
||||||
|
self.assertNotIn("switch to the author mcp session", action)
|
||||||
|
self.assertNotIn("switch to the reviewer mcp session", action)
|
||||||
|
self.assertIn("diagnostic defect", action)
|
||||||
|
|
||||||
|
def test_true_missing_permission_still_reports(self):
|
||||||
|
profile = _profile(
|
||||||
|
"prgs-author",
|
||||||
|
"author",
|
||||||
|
["gitea.read", "gitea.issue.create"],
|
||||||
|
["gitea.pr.merge"],
|
||||||
|
)
|
||||||
|
reviewer = _profile(
|
||||||
|
"prgs-reviewer",
|
||||||
|
"reviewer",
|
||||||
|
["gitea.read", "gitea.pr.merge", "gitea.pr.approve"],
|
||||||
|
[],
|
||||||
|
)
|
||||||
|
with patch.object(mcp_server, "get_profile", return_value=profile), patch.object(
|
||||||
|
mcp_server.gitea_config,
|
||||||
|
"load_config",
|
||||||
|
return_value=_config({"prgs-author": profile, "prgs-reviewer": reviewer}),
|
||||||
|
), patch.object(
|
||||||
|
mcp_server.gitea_config, "is_runtime_switching_enabled", return_value=True
|
||||||
|
):
|
||||||
|
report = mcp_server._permission_block_report("gitea.pr.merge")
|
||||||
|
self.assertFalse(report.get("diagnostic_defect"), report)
|
||||||
|
self.assertEqual(report.get("missing_permission"), "gitea.pr.merge")
|
||||||
|
self.assertIn("prgs-reviewer", report.get("matching_configured_profiles") or [])
|
||||||
|
|
||||||
|
|
||||||
|
class Issue897GateRefusalMatrix(unittest.TestCase):
|
||||||
|
def _stale_parity(self) -> dict:
|
||||||
|
return {
|
||||||
|
"in_parity": True,
|
||||||
|
"stale": False,
|
||||||
|
"restart_required": True,
|
||||||
|
"determinable": True,
|
||||||
|
"startup_head": SHA_START,
|
||||||
|
"current_head": SHA_START,
|
||||||
|
"daemon_start_head": SHA_START,
|
||||||
|
"local_head": SHA_START,
|
||||||
|
"live_remote_head": SHA_LIVE,
|
||||||
|
"live_known": True,
|
||||||
|
"live_stale": True,
|
||||||
|
"mutation_safe": False,
|
||||||
|
"reasons": [
|
||||||
|
f"live remote master is {SHA_LIVE[:12]} but the MCP server "
|
||||||
|
f"started at {SHA_START[:12]}; the daemon is stale relative "
|
||||||
|
"to live master -- restart/reconnect before mutating"
|
||||||
|
],
|
||||||
|
}
|
||||||
|
|
||||||
|
def test_stale_plus_permitted_op_all_roles(self):
|
||||||
|
for name, role, permitted_op, allowed, forbidden, _forbidden_op in ROLE_MATRIX:
|
||||||
|
with self.subTest(profile=name, op=permitted_op):
|
||||||
|
profile = _profile(name, role, allowed, forbidden)
|
||||||
|
parity = self._stale_parity()
|
||||||
|
with patch.object(mcp_server, "get_profile", return_value=profile), patch.object(
|
||||||
|
mcp_server, "_current_master_parity", return_value=parity
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_master_parity_block", return_value=list(parity["reasons"])
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_runtime_mode_block", return_value=[]
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_ensure_matching_profile", return_value=None
|
||||||
|
), patch.object(
|
||||||
|
mcp_server.session_ctx,
|
||||||
|
"mutation_context_audit_fields",
|
||||||
|
return_value={"session_profile": name},
|
||||||
|
):
|
||||||
|
blocked = mcp_server._profile_permission_block(permitted_op)
|
||||||
|
self.assertIsNotNone(blocked, name)
|
||||||
|
assert blocked is not None
|
||||||
|
self.assertEqual(
|
||||||
|
blocked.get("blocker_kind"),
|
||||||
|
"runtime_reconnect_required",
|
||||||
|
blocked,
|
||||||
|
)
|
||||||
|
self.assertNotIn("permission_report", blocked, blocked)
|
||||||
|
self.assertTrue(blocked.get("restart_required"), blocked)
|
||||||
|
self.assertEqual(blocked.get("startup_head"), SHA_START, blocked)
|
||||||
|
self.assertEqual(blocked.get("live_remote_head"), SHA_LIVE, blocked)
|
||||||
|
action = (blocked.get("exact_safe_next_action") or "").lower()
|
||||||
|
self.assertIn("reconnect", action)
|
||||||
|
self.assertNotIn("call gitea_activate_profile with", action)
|
||||||
|
self.assertNotIn("switch to the author mcp session", action)
|
||||||
|
self.assertNotIn("switch to the reviewer mcp session", action)
|
||||||
|
|
||||||
|
def test_fresh_plus_forbidden_op_all_roles(self):
|
||||||
|
for name, role, _permitted, allowed, forbidden, forbidden_op in ROLE_MATRIX:
|
||||||
|
with self.subTest(profile=name, op=forbidden_op):
|
||||||
|
profile = _profile(name, role, allowed, forbidden)
|
||||||
|
with patch.object(mcp_server, "get_profile", return_value=profile), patch.object(
|
||||||
|
mcp_server, "_master_parity_block", return_value=[]
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_runtime_mode_block", return_value=[]
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_ensure_matching_profile", return_value=None
|
||||||
|
), patch.object(
|
||||||
|
mcp_server.session_ctx,
|
||||||
|
"mutation_context_audit_fields",
|
||||||
|
return_value={"session_profile": name},
|
||||||
|
), patch.object(
|
||||||
|
mcp_server.gitea_config,
|
||||||
|
"load_config",
|
||||||
|
return_value=_config({name: profile}),
|
||||||
|
), patch.object(
|
||||||
|
mcp_server.gitea_config, "is_runtime_switching_enabled", return_value=False
|
||||||
|
):
|
||||||
|
blocked = mcp_server._profile_permission_block(forbidden_op)
|
||||||
|
self.assertIsNotNone(blocked, name)
|
||||||
|
assert blocked is not None
|
||||||
|
self.assertEqual(blocked.get("blocker_kind"), "permission_denied", blocked)
|
||||||
|
self.assertIn("permission_report", blocked, blocked)
|
||||||
|
report = blocked["permission_report"]
|
||||||
|
self.assertEqual(report.get("missing_permission"), forbidden_op, report)
|
||||||
|
self.assertFalse(report.get("diagnostic_defect"), report)
|
||||||
|
# No runtime reconnect fields for pure permission denial
|
||||||
|
self.assertNotEqual(
|
||||||
|
blocked.get("blocker_kind"), "runtime_reconnect_required"
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_stale_plus_forbidden_op_both_causes_separated(self):
|
||||||
|
for name, role, _permitted, allowed, forbidden, forbidden_op in ROLE_MATRIX:
|
||||||
|
with self.subTest(profile=name, op=forbidden_op):
|
||||||
|
profile = _profile(name, role, allowed, forbidden)
|
||||||
|
parity = self._stale_parity()
|
||||||
|
stale_reason = parity["reasons"][0]
|
||||||
|
with patch.object(mcp_server, "get_profile", return_value=profile), patch.object(
|
||||||
|
mcp_server, "_current_master_parity", return_value=parity
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_master_parity_block", return_value=[stale_reason]
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_runtime_mode_block", return_value=[]
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_ensure_matching_profile", return_value=None
|
||||||
|
), patch.object(
|
||||||
|
mcp_server.session_ctx,
|
||||||
|
"mutation_context_audit_fields",
|
||||||
|
return_value={"session_profile": name},
|
||||||
|
):
|
||||||
|
# Gate collects both classes; force permission reason too.
|
||||||
|
with patch.object(
|
||||||
|
mcp_server,
|
||||||
|
"_profile_operation_gate",
|
||||||
|
return_value=[
|
||||||
|
stale_reason,
|
||||||
|
f"profile is not allowed to {forbidden_op}",
|
||||||
|
],
|
||||||
|
):
|
||||||
|
blocked = mcp_server._profile_permission_block(forbidden_op)
|
||||||
|
self.assertIsNotNone(blocked)
|
||||||
|
assert blocked is not None
|
||||||
|
self.assertEqual(
|
||||||
|
blocked.get("blocker_kind"), "runtime_reconnect_required", blocked
|
||||||
|
)
|
||||||
|
self.assertNotIn("permission_report", blocked, blocked)
|
||||||
|
self.assertIn("permission_block_reasons", blocked, blocked)
|
||||||
|
self.assertIn("stale_runtime_reasons", blocked, blocked)
|
||||||
|
classes = blocked.get("gate_reason_classes") or {}
|
||||||
|
self.assertTrue(classes.get("stale_runtime"), classes)
|
||||||
|
self.assertTrue(classes.get("permission"), classes)
|
||||||
|
|
||||||
|
def test_runtime_mode_block_no_permission_report(self):
|
||||||
|
profile = _profile(
|
||||||
|
"prgs-author",
|
||||||
|
"author",
|
||||||
|
["gitea.read", "gitea.issue.create"],
|
||||||
|
[],
|
||||||
|
)
|
||||||
|
runtime_reason = (
|
||||||
|
"runtime mode is 'dev-test' and the mutation targets the "
|
||||||
|
"production repository; dev/test runtimes must not mutate real "
|
||||||
|
"issues or PRs (ADR: stable control runtime vs dev runtime)"
|
||||||
|
)
|
||||||
|
with patch.object(mcp_server, "get_profile", return_value=profile), patch.object(
|
||||||
|
mcp_server, "_master_parity_block", return_value=[]
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_runtime_mode_block", return_value=[runtime_reason]
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_ensure_matching_profile", return_value=None
|
||||||
|
), patch.object(
|
||||||
|
mcp_server.session_ctx,
|
||||||
|
"mutation_context_audit_fields",
|
||||||
|
return_value={"session_profile": "prgs-author"},
|
||||||
|
):
|
||||||
|
blocked = mcp_server._profile_permission_block("gitea.issue.create")
|
||||||
|
self.assertIsNotNone(blocked)
|
||||||
|
assert blocked is not None
|
||||||
|
self.assertEqual(blocked.get("blocker_kind"), "runtime_mode_blocked", blocked)
|
||||||
|
self.assertNotIn("permission_report", blocked, blocked)
|
||||||
|
action = (blocked.get("exact_safe_next_action") or "").lower()
|
||||||
|
self.assertNotIn("call gitea_activate_profile with", action)
|
||||||
|
self.assertIn("stable control runtime", action)
|
||||||
|
|
||||||
|
|
||||||
|
class Issue897CreateIssueRegression(unittest.TestCase):
|
||||||
|
def test_create_issue_stale_daemon_never_missing_issue_create(self):
|
||||||
|
"""Regression AC: stale prgs-author create_issue must not claim missing create."""
|
||||||
|
profile = _profile(
|
||||||
|
"prgs-author",
|
||||||
|
"author",
|
||||||
|
[
|
||||||
|
"gitea.read",
|
||||||
|
"gitea.issue.create",
|
||||||
|
"gitea.issue.comment",
|
||||||
|
"gitea.branch.create",
|
||||||
|
"gitea.branch.push",
|
||||||
|
"gitea.pr.create",
|
||||||
|
"gitea.pr.comment",
|
||||||
|
"gitea.repo.commit",
|
||||||
|
],
|
||||||
|
[],
|
||||||
|
)
|
||||||
|
stale_reason = (
|
||||||
|
f"live remote master is {SHA_LIVE[:12]} but the MCP server started "
|
||||||
|
f"at {SHA_START[:12]}; the daemon is stale relative to live master "
|
||||||
|
"-- restart/reconnect before mutating"
|
||||||
|
)
|
||||||
|
parity = {
|
||||||
|
"in_parity": True,
|
||||||
|
"stale": False,
|
||||||
|
"restart_required": True,
|
||||||
|
"determinable": True,
|
||||||
|
"startup_head": SHA_START,
|
||||||
|
"current_head": SHA_START,
|
||||||
|
"daemon_start_head": SHA_START,
|
||||||
|
"local_head": SHA_START,
|
||||||
|
"live_remote_head": SHA_LIVE,
|
||||||
|
"live_known": True,
|
||||||
|
"live_stale": True,
|
||||||
|
"mutation_safe": False,
|
||||||
|
"reasons": [stale_reason],
|
||||||
|
}
|
||||||
|
|
||||||
|
with patch.object(mcp_server, "get_profile", return_value=profile), patch.object(
|
||||||
|
mcp_server, "_current_master_parity", return_value=parity
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_master_parity_block", return_value=[stale_reason]
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_runtime_mode_block", return_value=[]
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_ensure_matching_profile", return_value=None
|
||||||
|
), patch.object(
|
||||||
|
mcp_server.session_ctx,
|
||||||
|
"mutation_context_audit_fields",
|
||||||
|
return_value={"session_profile": "prgs-author"},
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_mutation_config_authority_block", return_value=None
|
||||||
|
), patch.object(
|
||||||
|
mcp_server, "_session_context_mutation_block", return_value=None
|
||||||
|
):
|
||||||
|
blocked = mcp_server._profile_permission_block(
|
||||||
|
"gitea.issue.create", remote="prgs"
|
||||||
|
)
|
||||||
|
|
||||||
|
self.assertIsNotNone(blocked)
|
||||||
|
assert blocked is not None
|
||||||
|
self.assertEqual(blocked.get("blocker_kind"), "runtime_reconnect_required")
|
||||||
|
self.assertNotIn("permission_report", blocked)
|
||||||
|
# Even if a caller still built a raw report, holds-check must not claim missing.
|
||||||
|
with patch.object(mcp_server, "get_profile", return_value=profile):
|
||||||
|
raw = mcp_server._permission_block_report("gitea.issue.create")
|
||||||
|
self.assertIsNone(raw.get("missing_permission"), raw)
|
||||||
|
self.assertNotEqual(raw.get("missing_permission"), "gitea.issue.create")
|
||||||
|
|
||||||
|
def test_permission_report_for_gate_reasons_skips_stale(self):
|
||||||
|
stale = (
|
||||||
|
f"live remote master is {SHA_LIVE[:12]} but the MCP server started "
|
||||||
|
f"at {SHA_START[:12]}; the daemon is stale relative to live master "
|
||||||
|
"-- restart/reconnect before mutating"
|
||||||
|
)
|
||||||
|
self.assertIsNone(
|
||||||
|
mcp_server._permission_report_for_gate_reasons(
|
||||||
|
"gitea.issue.comment", [stale]
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
@@ -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.lease_views import render_leases_page
|
||||||
from webui.queue_loader import load_queue_snapshot, snapshot_to_dict as queue_snapshot_to_dict
|
from webui.queue_loader import load_queue_snapshot, snapshot_to_dict as queue_snapshot_to_dict
|
||||||
from webui.queue_views import render_queue_page
|
from webui.queue_views import render_queue_page
|
||||||
|
from webui.traffic_loader import load_traffic_snapshot, snapshot_to_dict as traffic_snapshot_to_dict
|
||||||
|
from webui.traffic_views import render_traffic_page
|
||||||
from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as worktree_snapshot_to_dict
|
from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as worktree_snapshot_to_dict
|
||||||
from webui.worktree_views import render_worktrees_page
|
from webui.worktree_views import render_worktrees_page
|
||||||
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
|
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
|
||||||
@@ -80,6 +82,7 @@ def _stub_page(title: str, description: str) -> HTMLResponse:
|
|||||||
|
|
||||||
|
|
||||||
_LEGACY_PAGES = (
|
_LEGACY_PAGES = (
|
||||||
|
("/traffic", "Traffic", "workflow traffic-control view (#640)"),
|
||||||
("/queue", "Queue", "live PR and issue dashboard (#429)"),
|
("/queue", "Queue", "live PR and issue dashboard (#429)"),
|
||||||
("/projects", "Projects", "registry and onboarding (#427)"),
|
("/projects", "Projects", "registry and onboarding (#427)"),
|
||||||
("/prompts", "Prompts", "canonical workflow prompt library (#428)"),
|
("/prompts", "Prompts", "canonical workflow prompt library (#428)"),
|
||||||
@@ -200,6 +203,15 @@ async def api_queue(_request: Request) -> JSONResponse:
|
|||||||
return JSONResponse(queue_snapshot_to_dict(load_queue_snapshot()))
|
return JSONResponse(queue_snapshot_to_dict(load_queue_snapshot()))
|
||||||
|
|
||||||
|
|
||||||
|
async def traffic(_request: Request) -> HTMLResponse:
|
||||||
|
snapshot = load_traffic_snapshot()
|
||||||
|
return HTMLResponse(render_traffic_page(snapshot))
|
||||||
|
|
||||||
|
|
||||||
|
async def api_traffic(_request: Request) -> JSONResponse:
|
||||||
|
return JSONResponse(traffic_snapshot_to_dict(load_traffic_snapshot()))
|
||||||
|
|
||||||
|
|
||||||
def _load_project_registry() -> tuple[ProjectRegistry | None, RegistryError | None]:
|
def _load_project_registry() -> tuple[ProjectRegistry | None, RegistryError | None]:
|
||||||
"""Load the registry, converting validation failure into a fail-closed pair."""
|
"""Load the registry, converting validation failure into a fail-closed pair."""
|
||||||
try:
|
try:
|
||||||
@@ -736,6 +748,8 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
|
|||||||
Route("/system-health", system_health, methods=["GET"]),
|
Route("/system-health", system_health, methods=["GET"]),
|
||||||
Route("/queue", queue, methods=["GET"]),
|
Route("/queue", queue, methods=["GET"]),
|
||||||
Route("/api/queue", api_queue, methods=["GET"]),
|
Route("/api/queue", api_queue, methods=["GET"]),
|
||||||
|
Route("/traffic", traffic, methods=["GET"]),
|
||||||
|
Route("/api/traffic", api_traffic, methods=["GET"]),
|
||||||
Route("/projects", projects, methods=["GET"]),
|
Route("/projects", projects, methods=["GET"]),
|
||||||
Route("/projects/{project_id}", project_detail, methods=["GET"]),
|
Route("/projects/{project_id}", project_detail, methods=["GET"]),
|
||||||
Route("/api/projects", api_projects, methods=["GET"]),
|
Route("/api/projects", api_projects, methods=["GET"]),
|
||||||
|
|||||||
@@ -182,10 +182,17 @@ def _extract_reviewer_leases(
|
|||||||
parsed = parse_reviewer_lease_comment(comment.get("body") or "")
|
parsed = parse_reviewer_lease_comment(comment.get("body") or "")
|
||||||
if not parsed:
|
if not parsed:
|
||||||
continue
|
continue
|
||||||
|
subject_pr = parsed.get("pr_number") or pr_number
|
||||||
leases.append(
|
leases.append(
|
||||||
{
|
{
|
||||||
**parsed,
|
**parsed,
|
||||||
"pr_number": parsed.get("pr_number") or pr_number,
|
"pr_number": subject_pr,
|
||||||
|
# The lease subject is the PR, never the linked issue: a
|
||||||
|
# reviewer lease on PR #N must not be attributed to issue #N
|
||||||
|
# or to the issue that PR closes (#640).
|
||||||
|
"kind": "pr",
|
||||||
|
"number": subject_pr,
|
||||||
|
"role": "reviewer",
|
||||||
"comment_id": comment.get("id"),
|
"comment_id": comment.get("id"),
|
||||||
"author": (comment.get("user") or {}).get("login"),
|
"author": (comment.get("user") or {}).get("login"),
|
||||||
"created_at": comment.get("created_at"),
|
"created_at": comment.get("created_at"),
|
||||||
|
|||||||
@@ -41,6 +41,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
|
|||||||
NavItem("/system-health", "System health"),
|
NavItem("/system-health", "System health"),
|
||||||
)),
|
)),
|
||||||
NavGroup("Traffic", (
|
NavGroup("Traffic", (
|
||||||
|
NavItem("/traffic", "Traffic control"),
|
||||||
NavItem("/queue", "Queue"),
|
NavItem("/queue", "Queue"),
|
||||||
NavItem("/leases", "Leases"),
|
NavItem("/leases", "Leases"),
|
||||||
NavItem("/actions", "Actions"),
|
NavItem("/actions", "Actions"),
|
||||||
|
|||||||
+34
-3
@@ -4,7 +4,7 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import os
|
import os
|
||||||
import re
|
import re
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass, field
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from typing import Any, Callable
|
from typing import Any, Callable
|
||||||
from urllib.parse import urlparse
|
from urllib.parse import urlparse
|
||||||
@@ -31,10 +31,20 @@ class PaginationMeta:
|
|||||||
|
|
||||||
@dataclass(frozen=True)
|
@dataclass(frozen=True)
|
||||||
class QueueItem:
|
class QueueItem:
|
||||||
|
"""One queue row.
|
||||||
|
|
||||||
|
``extra`` holds *display* strings for the queue page (values are truncated
|
||||||
|
or humanized for rendering). ``signals`` holds the *authoritative* typed
|
||||||
|
values taken straight from the Gitea payload, for consumers that classify
|
||||||
|
or pin state rather than render it (#640). Never derive identity or
|
||||||
|
concurrency decisions from ``extra``.
|
||||||
|
"""
|
||||||
|
|
||||||
number: int
|
number: int
|
||||||
title: str
|
title: str
|
||||||
badges: tuple[str, ...]
|
badges: tuple[str, ...]
|
||||||
extra: dict[str, str]
|
extra: dict[str, str]
|
||||||
|
signals: dict[str, Any] = field(default_factory=dict)
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True)
|
@dataclass(frozen=True)
|
||||||
@@ -134,21 +144,37 @@ def _format_pr_item(pr: dict, badges: tuple[str, ...]) -> QueueItem:
|
|||||||
"mergeable" if mergeable is True else "conflicted" if mergeable is False else "unknown"
|
"mergeable" if mergeable is True else "conflicted" if mergeable is False else "unknown"
|
||||||
)
|
)
|
||||||
linked = _extract_linked_issue(pr.get("title"), pr.get("body"))
|
linked = _extract_linked_issue(pr.get("title"), pr.get("body"))
|
||||||
|
head_sha = str(head.get("sha") or "")
|
||||||
|
labels = tuple(
|
||||||
|
str(lb.get("name") or "") for lb in (pr.get("labels") or []) if lb.get("name")
|
||||||
|
)
|
||||||
return QueueItem(
|
return QueueItem(
|
||||||
number=int(pr["number"]),
|
number=int(pr["number"]),
|
||||||
title=str(pr.get("title") or ""),
|
title=str(pr.get("title") or ""),
|
||||||
badges=badges,
|
badges=badges,
|
||||||
extra={
|
extra={
|
||||||
"branch": f"{head.get('ref', '?')} → {base.get('ref', '?')}",
|
"branch": f"{head.get('ref', '?')} → {base.get('ref', '?')}",
|
||||||
"head_sha": str(head.get("sha") or "")[:12],
|
# Display only — truncated. Pin against signals["head_sha"] instead.
|
||||||
|
"head_sha": head_sha[:12],
|
||||||
"mergeable": merge_label,
|
"mergeable": merge_label,
|
||||||
"linked_issue": str(linked) if linked is not None else "",
|
"linked_issue": str(linked) if linked is not None else "",
|
||||||
},
|
},
|
||||||
|
signals={
|
||||||
|
"head_sha": head_sha,
|
||||||
|
"head_ref": str(head.get("ref") or ""),
|
||||||
|
"base_ref": str(base.get("ref") or ""),
|
||||||
|
"mergeable": mergeable if isinstance(mergeable, bool) else None,
|
||||||
|
"labels": labels,
|
||||||
|
"linked_issue": linked,
|
||||||
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
|
def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
|
||||||
labels = ", ".join(lb.get("name", "") for lb in issue.get("labels", []))
|
label_names = tuple(
|
||||||
|
str(lb.get("name") or "") for lb in (issue.get("labels") or []) if lb.get("name")
|
||||||
|
)
|
||||||
|
labels = ", ".join(label_names)
|
||||||
assignee = (issue.get("assignee") or {}).get("login", "")
|
assignee = (issue.get("assignee") or {}).get("login", "")
|
||||||
return QueueItem(
|
return QueueItem(
|
||||||
number=int(issue["number"]),
|
number=int(issue["number"]),
|
||||||
@@ -159,6 +185,11 @@ def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
|
|||||||
"assignee": assignee or "unassigned",
|
"assignee": assignee or "unassigned",
|
||||||
"state": str(issue.get("state") or ""),
|
"state": str(issue.get("state") or ""),
|
||||||
},
|
},
|
||||||
|
signals={
|
||||||
|
"labels": label_names,
|
||||||
|
"assignee": assignee,
|
||||||
|
"state": str(issue.get("state") or ""),
|
||||||
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -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