Compare commits

...
Author SHA1 Message Date
sysadmin 6e6ca94338 fix(gate): stop classifying stale-runtime blocks as permission denials (Closes #897)
Stale-runtime and runtime-mode mutation refusals previously shared the
permission-denial channel, so permission_report claimed a missing op the
active profile already held and recommended gitea_activate_profile.
Typed blocker_kind payloads report reconnect-only recovery for staleness,
omit permission_report for non-permission gates, and fail closed when a
permission_report would invent a missing permission the profile holds.
2026-07-25 01:47:35 -04:00
sysadmin 2f4dec8323 Merge pull request 'feat(restart): pre-restart drain proof and hard gate (Closes #661)' (#882) from feat/issue-661-drain-proof-hard-gate into master 2026-07-24 23:40:02 -05:00
sysadmin 3a9d634c17 fix(drain-proof): bind acknowledgement coverage to session identity (#661)
Review 582 (REQUEST_CHANGES at 95178349) found a residual fail-open of the
same class the PR set out to close. Acknowledgement coverage was decided by
comparing a count against a count:

    covers_live_sessions = live_count_known and acked_count >= sessions_live_other

Nothing bound an acknowledgement to the identity of a session that actually
owed one, so acknowledgements supplied for the requesting session and for a
session that does not exist satisfied the obligations of two live sessions
that never answered - minting a clean, correctly signed proof and an allow
verdict from the restart gate.

Coverage is now derived from authoritative impact-report evidence:

- New `_required_ack_sessions()` derives the required session ids from the
  report itself, via `ack_state` keys and/or `affected_sessions` filtered on
  `live and not is_requester`. The requester is excluded only on explicit
  `is_requester` evidence, never inferred.
- When both views are present they must name the same set, and the result is
  reconciled against `counts.sessions_live_other`. Missing, malformed,
  duplicated, contradictory, or unreconcilable identity evidence fails closed
  and outranks every permitting path, including the timeout policy.
- Coverage requires every required id to carry an explicit acknowledgement
  token keyed by that id. Acknowledgements for the requester, for unknown
  ids, or for fabricated ids never increase coverage.
- Caller-supplied acknowledgement cardinality is no longer proof of anything.

Failure propagates unchanged through `acks_or_timeout` -> `proof.clean` ->
`failed_checks` -> `gate_apply_restart` verdict `deny` / `allow=False`.

The earlier missing-acknowledgement remediation is preserved in full: absent,
None, non-mapping, empty, partial, stale, and unparseable acks still fail
closed, `ack_timeout_policy_applied` stays strict `value is True`, and the
legitimate zero-live-sessions and explicit-timeout paths still pass.

Reviewer's reproduction, before and after this commit:

    sessions_live_other = 2
    report ack_state    = {'other-0': 'pending', 'other-1': 'pending'}
    supplied acks       = {'req': 'ack', 'totally-bogus-session': 'ack'}

    before: acks_or_timeout = True  | proof.clean = True  | gate allow
    after:  acks_or_timeout = False | proof.clean = False | gate deny

Tests: 22 new cases in `AcknowledgementIdentityBindingTests` covering the
reviewer's exact exploit, wrong-ids-with-sufficient-count, partial identity
match, requester-only acks, fabricated ids, unproven per-session states,
missing/malformed/contradictory identity evidence, count mismatch, and the
preserved success paths.

Verification:
- `pytest tests/test_drain_proof.py` -> 61 passed, 56 subtests
  (baseline at 95178349: 39 passed, 26 subtests)
- Restart surface (6 modules) -> 154 passed, 68 subtests, exit 0
  (baseline at 95178349: 132 passed, 38 subtests)
- Full `pytest tests/` -> 5177 passed vs baseline 5155 passed; the 23
  failures are identical in both runs and pre-exist at 95178349.

Scope: drain_proof.py, tests/test_drain_proof.py.

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01VRUZAf3Fr5n3kqhhiayN6C
(cherry picked from commit 4193b63f415b066ee292386c2c89bc3d2651a0cc)
2026-07-25 00:04:17 -04:00
jcwalker3 2068bae341 Merge branch 'master' into feat/issue-661-drain-proof-hard-gate 2026-07-24 22:34:27 -05:00
sysadmin 7af40fb5ff Merge pull request 'fix(allocator): exclude vision/roadmap/umbrella coordination containers (Closes #854)' (#883) from fix/issue-854-semantic-container-exclusion into master 2026-07-24 22:27:58 -05:00
sysadmin 9517834913 Merge commit '578c44b685a7ff5b01006c5e398bfac9863e0d8d' into feat/issue-661-drain-proof-hard-gate 2026-07-24 22:37:49 -04:00
sysadminandClaude Opus 5 824c42f7e3 fix(drain): fail closed on missing or unproven acknowledgement evidence (#661)
The acks_or_timeout check treated an absent `acks` key as proof that no
session needed to acknowledge: `drain_state.get("acks") or {}` collapsed
absent, None, and empty into the same value, and the resulting empty mapping
satisfied `no_sessions_to_ack`. The impact report's counts.sessions_live_other
was never consulted, so absence of evidence was read as evidence of absence.

Reproduced at head 1cbbde0089: with
sessions_live_other = 3 and the acknowledgement key absent, acks_or_timeout
passed with detail "no other live sessions required to acknowledge", the proof
minted clean, and gate_apply_restart returned verdict allow — a restart
authorized against three live sessions with zero acknowledgement evidence, and
the resulting artifact carried a valid signature.

Whether acknowledgement is required is now derived from the impact report,
never from the shape of the drain state:

- _live_session_count() reads counts.sessions_live_other and returns None for a
  missing, malformed, negative, or bool value, so an unreadable report fails
  closed instead of reading as "nobody was live".
- Absent, None, non-mapping, empty, partially-covering, and unparseable or
  stale acknowledgement data all fail closed while live sessions require
  acknowledgement.
- _is_acknowledged() no longer coerces with str(); only an explicit
  "ack"/"acked"/"acknowledged" string counts, so None, timestamps, and
  "pending"/"stale" markers are never read as an acknowledgement.
- Present-but-unacknowledged entries fail closed even when the report claims
  zero live sessions: that contradiction is not safe to resolve in favour of
  the restart.
- ack_timeout_policy_applied stays strict (`value is True`), so an absent, null,
  or non-boolean value cannot open the gate on its own.

The genuine no-other-live-sessions case still passes, now justified by the
report proving sessions_live_other == 0 rather than by the absence of data.

Adds AcknowledgementFailClosedTests: 14 cases / 26 subtests covering missing,
null, empty, malformed, stale, partial-coverage, and unproven-count inputs,
the valid-acknowledgement and zero-live-session paths, timeout-policy
strictness, and that a failed check blocks proof.clean, verification, and the
restart gate.

Restart-surface suite: 132 passed, 38 subtests (branch baseline 118 passed,
12 subtests; +14 new tests, no regressions).

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01VEaP3TohHLFWkp3Z2mmuZw
2026-07-24 22:36:25 -04:00
jcwalker3 578c44b685 Merge branch 'master' into feat/issue-661-drain-proof-hard-gate 2026-07-24 21:28:04 -05:00
jcwalker3 3b68d15593 Merge branch 'master' into fix/issue-854-semantic-container-exclusion 2026-07-24 21:27:55 -05:00
sysadmin a4c73766f4 Merge pull request 'feat(webui): Workflow traffic-control view (Phase 1) (Closes #640)' (#885) from issue-640 into master 2026-07-24 21:10:51 -05:00
jcwalker3 9f686253eb Merge branch 'master' into feat/issue-661-drain-proof-hard-gate 2026-07-24 21:06:49 -05:00
jcwalker3 b2e28428a4 Merge branch 'master' into fix/issue-854-semantic-container-exclusion 2026-07-24 21:06:41 -05:00
jcwalker3 95e4aae287 Merge branch 'master' into issue-640 2026-07-24 21:06:33 -05:00
sysadmin dac40ab9b3 docs(webui): update traffic state vocabulary docs and app nav for #640 2026-07-24 21:33:29 -04:00
sysadmin ccde9e8f11 Merge pull request 'fix(webui): migrate Starlette TestClient to httpx2 (Closes #682)' (#884) from fix/issue-682-starlette-httpx2 into master 2026-07-24 18:11:02 -05:00
sysadminandClaude Opus 4.8 1948d3dc21 fix(webui): repair traffic live path contracts for #640 review
Address PR #885 REQUEST_CHANGES: full head_sha pins from queue signals,
reviewer leases keyed by pr_number only, claim inventory via entries,
live-path fixture tests, and traffic state vocabulary docs.

Closes #640 (re-review at new head)

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-24 18:31:28 -04:00
sysadmin 069a9af7e6 feat(webui): implement workflow traffic-control view (Closes #640) 2026-07-24 17:34:50 -04:00
sysadmin 1cbbde0089 feat(restart): pre-restart drain proof and hard gate (#661)
Add `drain_proof.py`: a machine-verifiable DrainProof artifact plus a
fail-closed verifier and the hard gate the sanctioned restart-apply path
must consult, so a restart can never proceed on a stale or false "ready"
claim (#655 umbrella, child of #658 coordinator / #659 drain / #660
checkpoints).

- DrainProof: HMAC-SHA256 keyed proof-id over a canonical body using a
  per-process secret -> non-forgeable within the process; a proof minted in
  a prior daemon process will not verify after restart. Short TTL (120s).
- build_drain_proof(): mints the proof from the #658 impact report + the
  drain-mode outcomes. Checklist: no in-flight mutations, assignments
  stopped, checkpoints complete, handoffs ok, leases handled, acks-or-
  timeout. Every check fails closed on missing/ambiguous evidence; the
  no-in-flight-mutations and leases-handled checks are derived from the
  authoritative impact report, not self-reported.
- verify_drain_proof(): fail-closed — rejects missing, malformed, expired,
  signature-mismatched (forged/tampered/prior-process), unclean, or
  stale-fingerprint proofs; recomputes cleanliness from the checks rather
  than trusting the flag.
- gate_apply_restart(): allow only on a valid clean proof; deny -> durable
  incident descriptor; break-glass is the only bypass and is never silent.
- Checkpoint completeness is a supplied input, not a hard dependency on the
  (still-unmerged #660) checkpoint schema.

Wire the gate into gitea_request_mcp_restart: dry_run=False now enforces the
hard gate (drain_proof_json required; break-glass via request_break_glass +
GITEA_BREAKGLASS_RESTART_AUTHORIZATION env). The tool still performs no
actual restart — execution remains a further child.

Tests: tests/test_drain_proof.py — 25 cases covering AC#1-4 (apply without
proof denied, successful drain verifiable, open unsafe mutation fails,
pass/fail/expired), forgery/tamper/wrong-secret/stale-fingerprint rejection,
break-glass bypass, and secret hygiene. 25/25 pass (coordinator suite
unaffected: 40/40 together).

Links #652 #653 #655 #658 #659 #660.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01E7Fv9Bp2XWgvaWa4M1kdR7
(cherry picked from commit e7bcc952bb3e820fda95acbecefeaebfa5f8fcff)
2026-07-24 17:17:13 -04:00
sysadmin 0a78da39e5 fix(allocator): exclude vision/roadmap/umbrella coordination containers (#854)
#844 only caught epic-shaped child-only records. Live allocation still
selected product vision (#652), phased roadmap (#653), and umbrella (#655)
as implement targets. Extend pre-rank semantic classification with body
markers and container labels for those coordination records, keep title-
only and incidental mentions eligible, and add a live-equivalent canary.

Closes #854
2026-07-24 17:17:05 -04:00
14 changed files with 4303 additions and 99 deletions
+81 -17
View File
@@ -133,10 +133,15 @@ ROLE_ACTIONS: dict[str, tuple[tuple[str, ...], tuple[str, ...]]] = {
# Body phrases that prove an issue is an implementation container, not a
# unit of direct author work (#844). Matched case-insensitively against the
# issue body. Title alone is never sufficient (ordinary issues may mention
# "epic" incidentally).
# unit of direct author work (#844 / #854). Matched case-insensitively against
# the issue body. Title alone is never sufficient (ordinary issues may mention
# "epic", "roadmap", "vision", or "umbrella" incidentally).
#
# #854 extends the #844 marker set so product-vision (#652), phased-roadmap
# (#653), and umbrella (#655) coordination records — which do not use the word
# "epic" — are classified with the same semantic exclusion as epic containers.
_CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
# Epic / child-only (#844, live #631)
"implementation is delivered via child issues only",
"implementation is delivered through child issues only",
"implementation is delivered via child issues",
@@ -150,9 +155,26 @@ _CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
"coordination container",
"child-only container",
"implementation is delegated to child",
# Vision / roadmap / umbrella coordination (#854, live #652/#653/#655).
# Prefer authoritative non-implementation / child-only scope language over
# bare words like "roadmap" so ordinary implementable issues that mention
# a parent vision or roadmap stay eligible.
"do not implement features on this issue",
"implementing features on this roadmap issue",
"implementation is via linked children only",
"no product feature claimed complete on this issue alone",
"phased delivery roadmap and epic sequencing",
"this issue is the enduring source of truth",
"enduring source of truth for the",
"canonical product vision — enduring source of truth",
"canonical product vision - enduring source of truth",
"state: vision-active",
"state: roadmap-active",
)
# Explicit epic / umbrella labels (structured evidence preferred over title).
# Explicit epic / umbrella / vision / roadmap labels (structured evidence
# preferred over title). Tracker alone is *not* included — ordinary issues
# may carry a tracker label without being non-implementable containers.
_EPIC_LABELS: frozenset[str] = frozenset(
{
"type:epic",
@@ -161,6 +183,16 @@ _EPIC_LABELS: frozenset[str] = frozenset(
"scope:epic",
"type:umbrella",
"umbrella",
"kind:umbrella",
"scope:umbrella",
"type:vision",
"vision",
"kind:vision",
"scope:vision",
"type:roadmap",
"roadmap",
"kind:roadmap",
"scope:roadmap",
}
)
@@ -222,16 +254,45 @@ class WorkCandidate:
}
def _title_container_prefix(title_l: str) -> str | None:
"""Return a coordination-title prefix token if *title_l* uses one (#854).
Title prefixes alone never exclude; they only corroborate body/label
evidence. Ordinary issues may say "roadmap" or "vision" mid-title.
"""
for prefix, token in (
("epic:", "title_epic_prefix"),
("epic ", "title_epic_prefix"),
("umbrella:", "title_umbrella_prefix"),
("umbrella ", "title_umbrella_prefix"),
("roadmap:", "title_roadmap_prefix"),
("roadmap ", "title_roadmap_prefix"),
("product vision:", "title_vision_prefix"),
("product vision ", "title_vision_prefix"),
("vision:", "title_vision_prefix"),
("vision ", "title_vision_prefix"),
):
if title_l.startswith(prefix):
return token
return None
def classify_epic_or_child_only_container(
c: WorkCandidate,
) -> tuple[bool, str | None]:
"""Return whether *c* is an epic / child-only implementation container (#844).
"""Return whether *c* is a non-implementable coordination container (#844/#854).
Exclusion uses structured evidence first (labels, body scope language).
A bare title containing the word "epic" is **not** enough — ordinary
implementable issues may mention epics incidentally. A title that is
explicitly prefixed ``Epic:`` only counts when the body also proves
child-only / no-direct-implementation scope (or an epic label is present).
A bare title containing the words "epic", "roadmap", "vision", or
"umbrella" is **not** enough — ordinary implementable issues may mention
those terms incidentally. Explicit title prefixes (``Epic:``, ``Roadmap:``,
``Product vision:``, ``Umbrella:``) only count when the body also proves
child-only / no-direct-implementation scope (or a container label is
present).
Covers epic, product-vision, phased-roadmap, umbrella, and child-only
records so the allocator never assigns coordination containers as direct
author work.
PRs are never classified as containers here (they already have a head).
"""
@@ -245,25 +306,28 @@ def classify_epic_or_child_only_container(
title_l = title.lower()
body_hits = [m for m in _CHILD_ONLY_BODY_MARKERS if m in body_l]
title_epic_prefix = title_l.startswith("epic:") or title_l.startswith("epic ")
title_prefix = _title_container_prefix(title_l)
if epic_label:
detail = f"label={epic_label[0]}"
if body_hits:
detail = f"{detail}; body_marker={body_hits[0]!r}"
if title_prefix:
detail = f"{title_prefix}; {detail}"
return True, detail
if body_hits:
# Body proves child-only / umbrella scope. Title "Epic:" is corroborating
# but not required — containers without the word still exclude.
# Body proves child-only / vision / roadmap / umbrella scope. Title
# prefixes are corroborating but not required — containers without the
# title word still exclude.
detail = f"body_marker={body_hits[0]!r}"
if title_epic_prefix:
detail = f"title_epic_prefix; {detail}"
if title_prefix:
detail = f"{title_prefix}; {detail}"
return True, detail
# Title-only "Epic:" without body scope evidence is insufficient (#844 AC:
# eligibility does not rely solely on the word "Epic" in a title).
# Similarly, incidental "epic" mid-title without markers stays eligible.
# Title-only coordination prefix without body scope evidence is
# insufficient (#844/#854 AC: eligibility does not rely solely on a title
# word). Incidental mid-title mentions without markers stay eligible.
return False, None
+31
View File
@@ -57,6 +57,8 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| `/system-health` | System-health dashboard — readiness, version/uptime, dependencies, MCP namespaces, stale-runtime parity (#639) |
| `/queue` | Live PR and issue queue dashboard (#429) |
| `/api/queue` | JSON queue export with pagination metadata |
| `/traffic` | Workflow traffic-control view — runnable, leased, blocked, needs-controller, terminal-complete (#640) |
| `/api/traffic` | JSON traffic-control export with state classifications and next safe role actions |
| `/projects` | Project registry list with status and onboarding progress (#427, #635) |
| `/projects/{id}` | Project detail + onboarding checklist |
| `/api/v1/projects` | Versioned JSON registry export (#635) |
@@ -85,6 +87,35 @@ Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
local validator preview only (no Gitea mutations, no server-side storage).
### Traffic-control state vocabulary (#640)
The traffic view classifies each open issue/PR into exactly one bucket:
| Bucket | Meaning | Operator implication |
|--------|---------|----------------------|
| **runnable** | No active lease, no block reason, safe for its expected role | Next role may start work |
| **leased** | Active author claim or reviewer PR lease | Do not stomp; wait or adopt via role tools |
| **blocked** | Dependency, missing head pin, conflict, or unmet dependency | Author remediation first |
| **needs_controller** | Contaminated, controller-only diagnosis, or `status:blocked` | Controller only |
| **terminal_complete** | Reconciler / terminal-lock territory | Reconciler cleanup path |
`status:blocked` items route to **needs_controller**, not **blocked**:
`expected_role_for_candidate` sends them to the controller, and the blocker
reason renders in either bucket.
**Live path contracts (do not invent):**
- PR head pins come from `QueueItem.signals["head_sha"]` (full SHA). Display
`extra["head_sha"]` is truncated and must never be used for routing.
- Reviewer leases are keyed as `(pr, pr_number)` only — never via a linked
`issue_number` on the same lease marker.
- Issue claims come from `claim_inventory["entries"]`
(`issue_claim_heartbeat.build_claim_inventory`). There is no `active_claims`
key.
- Queue display badges are only: `blocked`, `claimed`, `duplicate`, `stale`,
`in-review`, `open`. Review verdicts (`request-changes`, `approved`) are
**not** queue badges; traffic does not invent them from the queue loader.
## System health API (#634)
`GET /api/v1/system/health` is the structured, read-only health surface for
+1015
View File
File diff suppressed because it is too large Load Diff
+354 -78
View File
@@ -2068,6 +2068,7 @@ import lease_lifecycle # noqa: E402
import lease_policy # noqa: E402
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
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 sentry_observability # noqa: E402 (#606 optional Sentry observability)
import sentry_incident_bridge # noqa: E402 (#607 Sentry→Gitea incident bridge)
@@ -9551,15 +9552,13 @@ def gitea_edit_pr(
if closing:
gate_reasons = _profile_operation_gate("gitea.pr.close")
if gate_reasons:
return {
"success": False,
"performed": False,
"pr_number": pr_number,
"requested_state": "closed",
"required_permission": "gitea.pr.close",
"reasons": gate_reasons,
"permission_report": _permission_block_report("gitea.pr.close"),
}
return _build_operation_gate_refusal(
"gitea.pr.close",
gate_reasons,
pr_number=pr_number,
requested_state="closed",
required_permission="gitea.pr.close",
)
h, o, r = _resolve(remote, host, org, repo)
auth = _auth(h)
@@ -13822,13 +13821,17 @@ def gitea_view_issue(
def _permission_block_report(required_operation: str,
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
refusal and never widens any permission, performs network I/O, or
raises (fail-soft: degrades to a minimal fail-closed report). Names
configured profiles only never auth references, tokens, endpoint
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 = {
"requested_operation": required_operation,
@@ -13840,6 +13843,7 @@ def _permission_block_report(required_operation: str,
"matching_configured_profiles": [],
"runtime_switching_supported": False,
"different_mcp_namespace_required": True,
"diagnostic_defect": False,
"exact_safe_next_action": (
"Ask the operator to fix GITEA_MCP_CONFIG/GITEA_MCP_PROFILE; "
"the active profile could not be resolved (fail closed)."),
@@ -13854,6 +13858,32 @@ def _permission_block_report(required_operation: str,
report["active_allowed_operations"] = (
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 = []
try:
config = gitea_config.load_config() or {}
@@ -13902,6 +13932,205 @@ def _permission_block_report(required_operation: str,
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:
# Normalize op first
try:
@@ -14078,7 +14307,7 @@ def _master_parity_block(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.*
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
not the promoted stable control runtime (#615) -- a dev/test or unknown
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)
if stale_reasons:
return stale_reasons
runtime_reasons = _runtime_mode_block(op)
if runtime_reasons:
return runtime_reasons
reasons: list[str] = []
reasons.extend(_master_parity_block(op))
reasons.extend(_runtime_mode_block(op))
try:
profile = get_profile()
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, profile["allowed_operations"], profile["forbidden_operations"])
if op_ok:
return []
return reasons
if _try_auto_switch_for_operation(op):
try:
@@ -14113,17 +14347,26 @@ def _profile_operation_gate(op: str) -> list[str]:
op_ok, op_reason = gitea_config.check_operation(
op, profile["allowed_operations"], profile["forbidden_operations"])
if op_ok:
return []
return reasons
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":
return ["profile has no configured allowed operations (fail closed)"]
if op_reason == "forbidden":
return [f"profile forbids '{op}'"]
if op_reason == "invalid-forbidden-entry":
return ["profile has an unrecognized forbidden operation entry (fail closed)"]
return [f"profile is not allowed to {op}"]
reasons.append(
"profile has no configured allowed operations (fail closed)"
)
elif op_reason == "forbidden":
reasons.append(f"profile forbids '{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:
@@ -14343,10 +14586,14 @@ def _session_context_mutation_block(
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*,
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 (
"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)
if reasons:
blocked = {
"success": False,
"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
return _build_operation_gate_refusal(
required_operation, reasons, **extra_fields
)
auth_block = _mutation_config_authority_block(required_operation)
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)."""
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"acquired": False,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
return _build_operation_gate_refusal(
"gitea.read", read_block, acquired=False
)
comment_block = _profile_operation_gate("gitea.pr.comment")
if comment_block:
return {
"success": False,
"acquired": False,
"reasons": comment_block,
"permission_report": _permission_block_report("gitea.pr.comment"),
}
return _build_operation_gate_refusal(
"gitea.pr.comment", comment_block, acquired=False
)
# task=acquire_reviewer_pr_lease so verify_preflight_purity runs shared #604
# 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")
if read_block:
return {
"success": False,
"acquired": False,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
return _build_operation_gate_refusal(
"gitea.read", read_block, acquired=False
)
comment_block = _profile_operation_gate("gitea.pr.comment")
if comment_block:
return {
"success": False,
"acquired": False,
"reasons": comment_block,
"permission_report": _permission_block_report("gitea.pr.comment"),
}
return _build_operation_gate_refusal(
"gitea.pr.comment", comment_block, acquired=False
)
merge_block = _profile_operation_gate("gitea.pr.merge")
if merge_block:
return {
@@ -19374,14 +19603,12 @@ def gitea_update_pr_branch_by_merge(
# Permission: author branch push / PR mutation surface.
push_block = _profile_operation_gate("gitea.branch.push")
if push_block:
return {
"success": False,
"performed": False,
"mutation_allowed": False,
"reasons": push_block,
"permission_report": _permission_block_report("gitea.branch.push"),
"role_kind": role,
}
return _build_operation_gate_refusal(
"gitea.branch.push",
push_block,
mutation_allowed=False,
role_kind=role,
)
if role != "author":
pre = pr_sync_status.assess_update_pr_branch_preflight(
@@ -22342,6 +22569,8 @@ def gitea_request_mcp_restart(
request_override: bool = False,
session_id: str | None = None,
limit: int = 200,
drain_proof_json: str | None = None,
request_break_glass: bool = False,
) -> dict:
"""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
would disrupt *before* any concurrent LLM work is destroyed.
This tool is **dry-run and never restarts anything.** The mutative apply
path is a separate child gated by a drain proof (non-goal here); calling
with ``dry_run=False`` still performs no restart and reports that apply is
not yet available.
This tool **never restarts a process.** In dry-run (the default) it returns
only the impact preview. With ``dry_run=False`` it enforces the #661 hard
gate: the apply request must present a valid, unexpired, clean drain proof
(``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
(``GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION``), never self-asserted by
@@ -22473,12 +22707,54 @@ def gitea_request_mcp_restart(
payload["requesting_session_id"] = sid
payload["operator_override_requested"] = bool(request_override)
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
if not dry_run:
payload["reasons"] = list(payload.get("reasons") or []) + [
"apply requested but not supported: sanctioned restart apply is "
"gated by a drain proof (separate child); no restart performed (#658)"
]
proof_obj: dict | None = None
proof_parse_error: str | None = None
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
+907
View File
@@ -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()
+408
View File
@@ -0,0 +1,408 @@
"""Tests for web UI workflow traffic-control view (#640)."""
import sys
import unittest
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient
from webui.app import create_app
from webui.traffic_loader import (
TrafficItem,
TrafficSnapshot,
load_traffic_snapshot,
snapshot_to_dict,
)
from webui.traffic_views import render_traffic_page
from allocator_service import WorkCandidate
class TestTrafficClassification(unittest.TestCase):
def test_runnable_candidate_classification(self):
cand = WorkCandidate(
kind="issue",
number=640,
state="open",
labels=("status:ready",),
title="Web Console: Workflow traffic-control view (Phase 1)",
priority=20,
)
snap = load_traffic_snapshot(candidates=[cand])
self.assertEqual(len(snap.runnable), 1)
self.assertEqual(snap.runnable[0].number, 640)
self.assertTrue(snap.runnable[0].is_safe)
self.assertEqual(snap.runnable[0].traffic_state, "runnable")
def test_blocked_dependency_candidate_classification(self):
cand = WorkCandidate(
kind="issue",
number=643,
state="open",
labels=("status:ready",),
title="Web Console: Requests & intent preview (Phase 2)",
priority=20,
dependency_unmet=True,
dependency_reason="issue#643 depends on unresolved issue(s) #640; they are not closed",
)
snap = load_traffic_snapshot(candidates=[cand])
self.assertEqual(len(snap.blocked), 1)
self.assertEqual(snap.blocked[0].number, 643)
self.assertFalse(snap.blocked[0].is_safe)
self.assertEqual(snap.blocked[0].traffic_state, "blocked")
self.assertIn("depends on unresolved issue(s) #640", snap.blocked[0].block_reason)
def test_leased_candidate_classification(self):
cand = WorkCandidate(
kind="issue",
number=640,
state="open",
labels=("status:in-progress",),
title="Web Console: Workflow traffic-control view (Phase 1)",
priority=20,
)
lease = {
"kind": "issue",
"number": 640,
"session_id": "prgs-author-12345",
"role": "author",
"status": "active",
}
snap = load_traffic_snapshot(candidates=[cand], leases=[lease])
self.assertEqual(len(snap.leased), 1)
self.assertEqual(snap.leased[0].number, 640)
self.assertEqual(snap.leased[0].traffic_state, "leased")
self.assertIsNotNone(snap.leased[0].lease_info)
def test_needs_controller_candidate_classification(self):
cand = WorkCandidate(
kind="issue",
number=700,
state="open",
labels=("status:blocked",),
title="Controller intervention needed",
priority=10,
blocked=True,
)
snap = load_traffic_snapshot(candidates=[cand])
self.assertEqual(len(snap.needs_controller), 1)
self.assertEqual(snap.needs_controller[0].number, 700)
class TestTrafficLoader(unittest.TestCase):
def test_snapshot_to_dict_export(self):
cand = WorkCandidate(
kind="issue",
number=640,
state="open",
labels=("status:ready",),
title="Traffic control test",
priority=20,
)
snap = load_traffic_snapshot(candidates=[cand])
data = snapshot_to_dict(snap)
self.assertEqual(data["project_id"], "gitea-tools")
self.assertEqual(len(data["runnable"]), 1)
self.assertTrue(data["inventory_complete"])
def test_fail_closed_error_handling(self):
with mock.patch("webui.traffic_loader.load_queue_snapshot", side_effect=RuntimeError("Gitea connection failed")):
snap = load_traffic_snapshot()
self.assertIsNotNone(snap.fetch_error)
self.assertIn("Failed to load traffic state", snap.fetch_error)
self.assertEqual(len(snap.runnable), 0)
self.assertFalse(snap.inventory_complete)
class TestTrafficLivePath(unittest.TestCase):
"""Live path tests: inject QueueSnapshot + LeaseSnapshot (no candidates=).
Covers the production ``load_traffic_snapshot()`` branch that ``/traffic``
and ``/api/traffic`` actually execute (#640 B1B5).
"""
FULL_SHA = "069a9af7e6aa2c2994e07199d1b0814819457017"
def _queue(
self,
*,
prs=(),
issues=(),
):
from webui.queue_loader import QueueSnapshot
return QueueSnapshot(
project_id="gitea-tools",
repo_label="Scaled-Tech-Consulting/Gitea-Tools",
prs=tuple(prs),
issues=tuple(issues),
pr_pagination=None,
issue_pagination=None,
fetch_error=None,
)
def _lease(
self,
*,
claim_inventory=None,
reviewer_leases=(),
):
from webui.lease_loader import LeaseSnapshot
return LeaseSnapshot(
project_id="gitea-tools",
repo_label="Scaled-Tech-Consulting/Gitea-Tools",
issue_lock=None,
claim_inventory=claim_inventory or {"entries": [], "counts": {}},
reviewer_leases=tuple(reviewer_leases),
duplicate_prs=(),
duplicate_branches=(),
collision_history=(),
fetch_error=None,
)
def test_live_pr_uses_full_head_sha_and_is_runnable(self):
from webui.queue_loader import QueueItem
pr = QueueItem(
number=885,
title="traffic control",
badges=("in-review",),
extra={"head_sha": self.FULL_SHA[:12], "linked_issue": "640"},
signals={
"head_sha": self.FULL_SHA,
"mergeable": True,
"labels": (),
"linked_issue": 640,
},
)
q = self._queue(prs=[pr])
l = self._lease()
snap = load_traffic_snapshot(
fetch_queue_snapshot=lambda: q,
fetch_lease_snapshot=lambda: l,
)
self.assertIsNone(snap.fetch_error)
self.assertEqual(len(snap.runnable), 1)
item = snap.runnable[0]
self.assertEqual(item.kind, "pr")
self.assertEqual(item.number, 885)
self.assertEqual(item.head_sha, self.FULL_SHA)
self.assertNotEqual(item.head_sha, self.FULL_SHA[:12])
self.assertIsNone(item.block_reason)
self.assertEqual(len(snap.blocked), 0)
def test_live_pr_without_head_sha_is_blocked(self):
from webui.queue_loader import QueueItem
pr = QueueItem(
number=1,
title="missing pin",
badges=("open",),
extra={"head_sha": ""},
signals={"head_sha": "", "mergeable": True, "labels": ()},
)
snap = load_traffic_snapshot(
fetch_queue_snapshot=lambda: self._queue(prs=[pr]),
fetch_lease_snapshot=lambda: self._lease(),
)
self.assertEqual(len(snap.blocked) + len(snap.needs_controller), 1)
item = (snap.blocked or snap.needs_controller)[0]
self.assertIn("missing head_sha", (item.block_reason or "").lower())
def test_reviewer_lease_keys_by_pr_not_linked_issue(self):
from webui.queue_loader import QueueItem
pr = QueueItem(
number=885,
title="leased pr",
badges=("in-review",),
extra={"head_sha": self.FULL_SHA[:12]},
signals={"head_sha": self.FULL_SHA, "mergeable": True, "labels": ()},
)
issue = QueueItem(
number=640,
title="linked issue",
badges=("open",),
extra={},
signals={"labels": ()},
)
# Marker-shaped record: has both pr_number and issue_number; must
# attach to the PR only (B2).
reviewer_lease = {
"pr_number": 885,
"issue_number": 640,
"phase": "validating",
"reviewer_identity": "sysadmin",
"session_id": "review-sess-1",
}
snap = load_traffic_snapshot(
fetch_queue_snapshot=lambda: self._queue(prs=[pr], issues=[issue]),
fetch_lease_snapshot=lambda: self._lease(reviewer_leases=[reviewer_lease]),
)
leased_prs = [i for i in snap.leased if i.kind == "pr" and i.number == 885]
self.assertEqual(len(leased_prs), 1)
self.assertEqual(leased_prs[0].lease_info.get("pr_number"), 885)
# Issue 640 must not inherit the reviewer lease just because issue_number
# is present on the marker.
for item in list(snap.leased) + list(snap.runnable) + list(snap.blocked):
if item.kind == "issue" and item.number == 640:
self.assertIsNone(
item.lease_info,
"reviewer lease must not attach to linked issue #640",
)
break
else:
self.fail("expected issue #640 in traffic snapshot")
def test_claim_inventory_entries_key_marks_issue_leased(self):
from webui.queue_loader import QueueItem
issue = QueueItem(
number=640,
title="claimed issue",
badges=("claimed",),
extra={},
signals={"labels": ("status:in-progress",)},
)
inventory = {
"entries": [
{
"issue_number": 640,
"status": "active",
"latest_heartbeat": {"session_id": "author-sess-9"},
"reasons": ["claim has structured heartbeat proof"],
}
],
"counts": {"active": 1},
"in_progress_total": 1,
}
snap = load_traffic_snapshot(
fetch_queue_snapshot=lambda: self._queue(issues=[issue]),
fetch_lease_snapshot=lambda: self._lease(claim_inventory=inventory),
)
leased_issues = [i for i in snap.leased if i.kind == "issue" and i.number == 640]
self.assertEqual(len(leased_issues), 1)
self.assertEqual(leased_issues[0].traffic_state, "leased")
def test_active_claims_key_is_ignored(self):
"""B3 regression: fictional ``active_claims`` must not create lease_info."""
from webui.queue_loader import QueueItem
issue = QueueItem(
number=640,
title="open issue",
badges=("open",),
extra={},
signals={"labels": ()},
)
# Only the broken key — must NOT produce lease_info. Entries-less
# inventory is empty (entries is the real claim_inventory key).
inventory = {
"active_claims": [
{
"kind": "issue",
"number": 640,
"issue_number": 640,
"status": "active",
},
],
"counts": {},
}
snap = load_traffic_snapshot(
fetch_queue_snapshot=lambda: self._queue(issues=[issue]),
fetch_lease_snapshot=lambda: self._lease(claim_inventory=inventory),
)
items = [
i
for i in (
list(snap.runnable)
+ list(snap.leased)
+ list(snap.blocked)
+ list(snap.needs_controller)
)
if i.kind == "issue" and i.number == 640
]
self.assertEqual(len(items), 1)
self.assertIsNone(
items[0].lease_info,
"active_claims is not a real inventory key; entries-only",
)
class TestTrafficRoutesAndRendering(unittest.TestCase):
def setUp(self):
self.client = TestClient(create_app())
def test_traffic_html_page_rendering(self):
cand1 = WorkCandidate(
kind="issue",
number=640,
state="open",
labels=("status:ready",),
title="Traffic View Implementation",
priority=20,
)
cand2 = WorkCandidate(
kind="issue",
number=643,
state="open",
labels=("status:ready",),
title="Dependent Feature",
priority=20,
dependency_unmet=True,
dependency_reason="issue#643 depends on unresolved issue(s) #640; they are not closed",
)
snap = load_traffic_snapshot(candidates=[cand1, cand2])
with mock.patch("webui.app.load_traffic_snapshot", return_value=snap):
response = self.client.get("/traffic")
self.assertEqual(response.status_code, 200)
self.assertIn("Workflow Traffic Control", response.text)
self.assertIn("1. Runnable Lanes", response.text)
self.assertIn("3. Blocked Items", response.text)
self.assertIn("Traffic View Implementation", response.text)
self.assertIn("depends on unresolved issue(s) #640", response.text)
def test_api_traffic_json_route(self):
cand = WorkCandidate(
kind="issue",
number=640,
state="open",
labels=("status:ready",),
title="Traffic View API Test",
priority=20,
)
snap = load_traffic_snapshot(candidates=[cand])
with mock.patch("webui.app.load_traffic_snapshot", return_value=snap):
response = self.client.get("/api/traffic")
self.assertEqual(response.status_code, 200)
data = response.json()
self.assertEqual(data["project_id"], "gitea-tools")
self.assertEqual(len(data["runnable"]), 1)
self.assertEqual(data["runnable"][0]["number"], 640)
def test_render_traffic_fail_closed_page(self):
snap = TrafficSnapshot(
project_id="gitea-tools",
repo_label="Scaled-Tech-Consulting/Gitea-Tools",
runnable=(),
leased=(),
blocked=(),
needs_controller=(),
terminal_complete=(),
next_roles=(),
fetch_error="Gitea credentials unavailable for gitea.prgs.cc",
inventory_complete=False,
)
html = render_traffic_page(snap)
self.assertIn("Traffic data unavailable", html)
self.assertIn("Fail closed", html)
self.assertNotIn("1. Runnable Lanes", html)
if __name__ == "__main__":
unittest.main()
+14
View File
@@ -42,6 +42,8 @@ from webui.lease_loader import load_lease_snapshot, snapshot_to_dict as lease_sn
from webui.lease_views import render_leases_page
from webui.queue_loader import load_queue_snapshot, snapshot_to_dict as queue_snapshot_to_dict
from webui.queue_views import render_queue_page
from webui.traffic_loader import load_traffic_snapshot, snapshot_to_dict as traffic_snapshot_to_dict
from webui.traffic_views import render_traffic_page
from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as worktree_snapshot_to_dict
from webui.worktree_views import render_worktrees_page
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
@@ -80,6 +82,7 @@ def _stub_page(title: str, description: str) -> HTMLResponse:
_LEGACY_PAGES = (
("/traffic", "Traffic", "workflow traffic-control view (#640)"),
("/queue", "Queue", "live PR and issue dashboard (#429)"),
("/projects", "Projects", "registry and onboarding (#427)"),
("/prompts", "Prompts", "canonical workflow prompt library (#428)"),
@@ -200,6 +203,15 @@ async def api_queue(_request: Request) -> JSONResponse:
return JSONResponse(queue_snapshot_to_dict(load_queue_snapshot()))
async def traffic(_request: Request) -> HTMLResponse:
snapshot = load_traffic_snapshot()
return HTMLResponse(render_traffic_page(snapshot))
async def api_traffic(_request: Request) -> JSONResponse:
return JSONResponse(traffic_snapshot_to_dict(load_traffic_snapshot()))
def _load_project_registry() -> tuple[ProjectRegistry | None, RegistryError | None]:
"""Load the registry, converting validation failure into a fail-closed pair."""
try:
@@ -736,6 +748,8 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/system-health", system_health, methods=["GET"]),
Route("/queue", queue, methods=["GET"]),
Route("/api/queue", api_queue, methods=["GET"]),
Route("/traffic", traffic, methods=["GET"]),
Route("/api/traffic", api_traffic, methods=["GET"]),
Route("/projects", projects, methods=["GET"]),
Route("/projects/{project_id}", project_detail, methods=["GET"]),
Route("/api/projects", api_projects, methods=["GET"]),
+8 -1
View File
@@ -182,10 +182,17 @@ def _extract_reviewer_leases(
parsed = parse_reviewer_lease_comment(comment.get("body") or "")
if not parsed:
continue
subject_pr = parsed.get("pr_number") or pr_number
leases.append(
{
**parsed,
"pr_number": parsed.get("pr_number") or pr_number,
"pr_number": subject_pr,
# The lease subject is the PR, never the linked issue: a
# reviewer lease on PR #N must not be attributed to issue #N
# or to the issue that PR closes (#640).
"kind": "pr",
"number": subject_pr,
"role": "reviewer",
"comment_id": comment.get("id"),
"author": (comment.get("user") or {}).get("login"),
"created_at": comment.get("created_at"),
+1
View File
@@ -41,6 +41,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
NavItem("/system-health", "System health"),
)),
NavGroup("Traffic", (
NavItem("/traffic", "Traffic control"),
NavItem("/queue", "Queue"),
NavItem("/leases", "Leases"),
NavItem("/actions", "Actions"),
+34 -3
View File
@@ -4,7 +4,7 @@ from __future__ import annotations
import os
import re
from dataclasses import dataclass
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Callable
from urllib.parse import urlparse
@@ -31,10 +31,20 @@ class PaginationMeta:
@dataclass(frozen=True)
class QueueItem:
"""One queue row.
``extra`` holds *display* strings for the queue page (values are truncated
or humanized for rendering). ``signals`` holds the *authoritative* typed
values taken straight from the Gitea payload, for consumers that classify
or pin state rather than render it (#640). Never derive identity or
concurrency decisions from ``extra``.
"""
number: int
title: str
badges: tuple[str, ...]
extra: dict[str, str]
signals: dict[str, Any] = field(default_factory=dict)
@dataclass(frozen=True)
@@ -134,21 +144,37 @@ def _format_pr_item(pr: dict, badges: tuple[str, ...]) -> QueueItem:
"mergeable" if mergeable is True else "conflicted" if mergeable is False else "unknown"
)
linked = _extract_linked_issue(pr.get("title"), pr.get("body"))
head_sha = str(head.get("sha") or "")
labels = tuple(
str(lb.get("name") or "") for lb in (pr.get("labels") or []) if lb.get("name")
)
return QueueItem(
number=int(pr["number"]),
title=str(pr.get("title") or ""),
badges=badges,
extra={
"branch": f"{head.get('ref', '?')}{base.get('ref', '?')}",
"head_sha": str(head.get("sha") or "")[:12],
# Display only — truncated. Pin against signals["head_sha"] instead.
"head_sha": head_sha[:12],
"mergeable": merge_label,
"linked_issue": str(linked) if linked is not None else "",
},
signals={
"head_sha": head_sha,
"head_ref": str(head.get("ref") or ""),
"base_ref": str(base.get("ref") or ""),
"mergeable": mergeable if isinstance(mergeable, bool) else None,
"labels": labels,
"linked_issue": linked,
},
)
def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
labels = ", ".join(lb.get("name", "") for lb in issue.get("labels", []))
label_names = tuple(
str(lb.get("name") or "") for lb in (issue.get("labels") or []) if lb.get("name")
)
labels = ", ".join(label_names)
assignee = (issue.get("assignee") or {}).get("login", "")
return QueueItem(
number=int(issue["number"]),
@@ -159,6 +185,11 @@ def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
"assignee": assignee or "unassigned",
"state": str(issue.get("state") or ""),
},
signals={
"labels": label_names,
"assignee": assignee,
"state": str(issue.get("state") or ""),
},
)
+449
View File
@@ -0,0 +1,449 @@
"""Traffic-control view loader for Phase 1 operator web console (#640).
Combines queue snapshots, inventory leases, dependency graph classifications,
and workflow dashboard rules to deliver full traffic-control visibility:
runnable, leased (in-progress), blocked (dependency/lock), needs-controller,
and terminal-complete candidates.
"""
from __future__ import annotations
from dataclasses import dataclass
from typing import Any, Callable, Sequence
from webui.project_registry import find_project, load_registry
from webui.queue_loader import load_queue_snapshot, QueueSnapshot
from webui.lease_loader import load_lease_snapshot, LeaseSnapshot
from workflow_dashboard import (
DashboardSnapshot,
QueueEntry,
RoleNextAction,
build_workflow_dashboard,
DASHBOARD_ROLES,
)
from allocator_service import WorkCandidate
@dataclass(frozen=True)
class TrafficItem:
kind: str # "issue" or "pr"
number: int
title: str
traffic_state: str # "runnable", "leased", "blocked", "needs_controller", "terminal_complete"
expected_role: str
safe_for_roles: tuple[str, ...]
badges: tuple[str, ...]
block_reason: str | None = None
lease_info: dict[str, Any] | None = None
head_sha: str | None = None
@property
def is_safe(self) -> bool:
return self.block_reason is None and bool(self.safe_for_roles)
def as_dict(self) -> dict[str, Any]:
return {
"kind": self.kind,
"number": self.number,
"title": self.title,
"traffic_state": self.traffic_state,
"expected_role": self.expected_role,
"safe_for_roles": list(self.safe_for_roles),
"badges": list(self.badges),
"block_reason": self.block_reason,
"lease_info": self.lease_info,
"head_sha": self.head_sha,
"is_safe": self.is_safe,
}
@dataclass(frozen=True)
class TrafficSnapshot:
project_id: str
repo_label: str
runnable: tuple[TrafficItem, ...]
leased: tuple[TrafficItem, ...]
blocked: tuple[TrafficItem, ...]
needs_controller: tuple[TrafficItem, ...]
terminal_complete: tuple[TrafficItem, ...]
next_roles: tuple[dict[str, Any], ...]
fetch_error: str | None = None
inventory_complete: bool = True
def as_dict(self) -> dict[str, Any]:
return {
"project_id": self.project_id,
"repo_label": self.repo_label,
"runnable": [i.as_dict() for i in self.runnable],
"leased": [i.as_dict() for i in self.leased],
"blocked": [i.as_dict() for i in self.blocked],
"needs_controller": [i.as_dict() for i in self.needs_controller],
"terminal_complete": [i.as_dict() for i in self.terminal_complete],
"next_roles": list(self.next_roles),
"fetch_error": self.fetch_error,
"inventory_complete": self.inventory_complete,
}
def _classify_traffic_item(
entry: QueueEntry,
*,
lease_info: dict[str, Any] | None = None,
) -> TrafficItem:
"""Classify a QueueEntry into a TrafficItem with explicit traffic state."""
badges = list(entry.badges)
block_reason = entry.block_reason
expected_role = entry.expected_role
entry_is_safe = entry.block_reason is None and bool(entry.safe_for_roles)
# Lease state is checked first: an item that is both leased and blocked is
# reported as leased. That is safe by construction — a leased item is never
# placed in the runnable lane — and it keeps the operator's attention on the
# session that currently owns the work. The blocker text still renders.
if lease_info is not None or "in-progress" in badges or "claimed" in badges:
state = "leased"
elif expected_role == "reconciler" or "terminal-lock" in badges:
state = "terminal_complete"
elif expected_role == "controller" or "contaminated" in badges or "needs-controller" in badges:
state = "needs_controller"
elif (
block_reason is not None
or "blocked" in badges
or "dependency-unmet" in badges
or "blocked-by-terminal" in badges
or "status:blocked" in badges
):
state = "blocked"
elif entry_is_safe:
state = "runnable"
else:
state = "needs_controller"
return TrafficItem(
kind=entry.kind,
number=entry.number,
title=entry.title,
traffic_state=state,
expected_role=expected_role,
safe_for_roles=entry.safe_for_roles,
badges=tuple(badges),
block_reason=block_reason,
lease_info=lease_info,
head_sha=entry.head_sha,
)
# Claim statuses from ``issue_claim_heartbeat.build_claim_inventory`` that mean
# a live worker currently holds the issue. Everything else (``stale``,
# ``phantom``, ``reclaimable``, ``not_claimed``) is reported through the
# dashboard's stale-lease channel and is never rendered as an active lease.
_ACTIVE_CLAIM_STATUSES = frozenset({"active", "awaiting_review"})
# Statuses that positively mean "not an active lease" for any lease record.
_INACTIVE_LEASE_STATUSES = frozenset(
{"expired", "stale", "released", "moot", "reclaimable", "phantom", "not_claimed"}
)
def _candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate]:
"""Build allocator candidates from the queue loader's authoritative signals.
Display badges (``blocked``/``claimed``/``duplicate``/``stale``/
``in-review``/``open``) are rendering hints, not routing state, so nothing
here branches on them. Every routing field comes from
``QueueItem.signals`` the raw Gitea payload values.
The queue loader reads ``/pulls`` and ``/issues`` only; it never fetches
review verdicts. ``request_changes_current_head`` / ``approval_on_current_head``
are therefore left at their fail-safe ``False`` rather than being guessed
from badges: an unproven approval must never route a PR to the merger.
"""
candidates: list[WorkCandidate] = []
for pr in q_snap.prs:
signals = pr.signals or {}
head_sha = str(signals.get("head_sha") or "").strip()
mergeable = signals.get("mergeable")
labels = tuple(str(x) for x in (signals.get("labels") or ()))
candidates.append(
WorkCandidate(
kind="pr",
number=pr.number,
state="open",
labels=labels,
title=pr.title,
# Full 40-char SHA from head.sha — never the 12-char display value.
head_sha=head_sha or None,
priority=5,
mergeable=mergeable is True,
blocked=mergeable is False or "status:blocked" in labels,
)
)
for issue in q_snap.issues:
signals = issue.signals or {}
labels = tuple(str(x) for x in (signals.get("labels") or ()))
lowered = {label.lower() for label in labels}
candidates.append(
WorkCandidate(
kind="issue",
number=issue.number,
state="open",
labels=labels,
title=issue.title,
priority=20 if "status:ready" in lowered else 10,
blocked="status:blocked" in lowered,
# A live claim by another session is not this session's work.
already_claimed_elsewhere="status:in-progress" in lowered,
)
)
return candidates
def _claim_lease_records(inventory: dict[str, Any] | None) -> list[dict[str, Any]]:
"""Normalize ``build_claim_inventory`` entries into lease records.
The inventory contract is ``{"entries", "counts", "heartbeat_lease_minutes",
"reclaim_after_minutes", "in_progress_total"}``. Each entry is keyed by
``issue_number``; the subject kind is therefore always ``issue``.
"""
entries = (inventory or {}).get("entries") or ()
records: list[dict[str, Any]] = []
for entry in entries:
if not isinstance(entry, dict):
continue
number = entry.get("issue_number")
if number is None:
continue
try:
number_int = int(number)
except (TypeError, ValueError):
continue
heartbeat = entry.get("latest_heartbeat") or {}
record = dict(entry)
record.update(
{
"kind": "issue",
"number": number_int,
"role": "author",
"lease_source": "issue-claim-heartbeat",
}
)
if isinstance(heartbeat, dict):
if heartbeat.get("session_id") and not record.get("session_id"):
record["session_id"] = heartbeat.get("session_id")
if heartbeat.get("author") and not record.get("author"):
record["author"] = heartbeat.get("author")
records.append(record)
return records
def _lease_subject(lease: dict[str, Any]) -> tuple[str, int] | None:
"""Return the ``(kind, number)`` a lease record actually covers.
Fails closed: a record that does not identify exactly one subject is
dropped rather than attributed to a guessed work item (#640 — never invent
a lease, and never attach a PR lease to a same-numbered issue).
"""
kind = str(lease.get("kind") or lease.get("work_kind") or "").strip().lower()
pr_number = lease.get("pr_number")
issue_number = lease.get("issue_number")
if kind not in ("pr", "issue"):
if pr_number is not None and issue_number is None:
kind = "pr"
elif issue_number is not None and pr_number is None:
kind = "issue"
else:
return None
number = lease.get("number")
if number is None:
number = lease.get("work_number")
if number is None:
number = pr_number if kind == "pr" else issue_number
if number is None:
return None
try:
return kind, int(number)
except (TypeError, ValueError):
return None
def _is_active_lease(lease: dict[str, Any]) -> bool:
"""True when the record proves a worker currently holds the item."""
if lease.get("stale") or lease.get("expired"):
return False
status = str(lease.get("status") or lease.get("lease_status") or "").strip().lower()
if status in _INACTIVE_LEASE_STATUSES:
return False
if lease.get("lease_source") == "issue-claim-heartbeat":
return status in _ACTIVE_CLAIM_STATUSES
return True
def load_traffic_snapshot(
*,
candidates: Sequence[WorkCandidate] | None = None,
leases: Sequence[dict[str, Any]] | None = None,
terminal_pr: int | None = None,
fetch_queue_snapshot: Callable[[], QueueSnapshot] | None = None,
fetch_lease_snapshot: Callable[[], LeaseSnapshot] | None = None,
project_id: str = "gitea-tools",
) -> TrafficSnapshot:
"""Load and compute the traffic-control snapshot."""
try:
reg = load_registry()
proj = find_project(reg, project_id)
repo_label = proj.remote_repo if proj else "Scaled-Tech-Consulting/Gitea-Tools"
except Exception:
repo_label = "Scaled-Tech-Consulting/Gitea-Tools"
# Injected candidates path (pure unit testing)
if candidates is not None:
dashboard = build_workflow_dashboard(
candidates=candidates,
leases=leases,
terminal_pr=terminal_pr,
inventory_complete=True,
)
return _build_traffic_snapshot_from_dashboard(
project_id=project_id,
repo_label=repo_label,
dashboard=dashboard,
leases=leases or (),
)
# Live snapshot loading
q_loader = fetch_queue_snapshot or load_queue_snapshot
l_loader = fetch_lease_snapshot or load_lease_snapshot
try:
q_snap = q_loader()
l_snap = l_loader()
except Exception as exc: # noqa: BLE001
return TrafficSnapshot(
project_id=project_id,
repo_label=repo_label,
runnable=(),
leased=(),
blocked=(),
needs_controller=(),
terminal_complete=(),
next_roles=(),
fetch_error=f"Failed to load traffic state: {exc}",
inventory_complete=False,
)
if q_snap.fetch_error or l_snap.fetch_error:
err = q_snap.fetch_error or l_snap.fetch_error
return TrafficSnapshot(
project_id=project_id,
repo_label=repo_label,
runnable=(),
leased=(),
blocked=(),
needs_controller=(),
terminal_complete=(),
next_roles=(),
fetch_error=err,
inventory_complete=False,
)
candidate_list = _candidates_from_queue_snapshot(q_snap)
raw_leases: list[dict[str, Any]] = _claim_lease_records(l_snap.claim_inventory)
for r_lease in l_snap.reviewer_leases or ():
if not isinstance(r_lease, dict):
continue
# Always pin reviewer leases to the PR subject, even if a linked
# issue_number is present on the marker (#640 B2).
normalized = dict(r_lease)
subject = normalized.get("pr_number") or normalized.get("number")
if subject is None:
continue
try:
pr_num = int(subject)
except (TypeError, ValueError):
continue
normalized["kind"] = "pr"
normalized["number"] = pr_num
normalized["pr_number"] = pr_num
normalized.setdefault("role", "reviewer")
raw_leases.append(normalized)
dashboard = build_workflow_dashboard(
candidates=candidate_list,
leases=raw_leases,
inventory_complete=q_snap.pr_pagination.inventory_complete if q_snap.pr_pagination else True,
)
return _build_traffic_snapshot_from_dashboard(
project_id=project_id,
repo_label=repo_label,
dashboard=dashboard,
leases=raw_leases,
)
def _build_traffic_snapshot_from_dashboard(
*,
project_id: str,
repo_label: str,
dashboard: DashboardSnapshot,
leases: Sequence[dict[str, Any]],
) -> TrafficSnapshot:
"""Classify dashboard entries into the 5 traffic state buckets."""
all_entries = dashboard.open_prs + dashboard.open_issues
# Map each active lease onto the exact work item it covers. Records whose
# subject cannot be determined, and claims that are stale/phantom/
# reclaimable, are deliberately dropped instead of guessed.
lease_map: dict[tuple[str, int], dict[str, Any]] = {}
for lease in leases:
if not isinstance(lease, dict) or not _is_active_lease(lease):
continue
subject = _lease_subject(lease)
if subject is not None:
lease_map[subject] = lease
runnable: list[TrafficItem] = []
leased: list[TrafficItem] = []
blocked: list[TrafficItem] = []
needs_controller: list[TrafficItem] = []
terminal_complete: list[TrafficItem] = []
for entry in all_entries:
l_info = lease_map.get((entry.kind, entry.number))
item = _classify_traffic_item(entry, lease_info=l_info)
if item.traffic_state == "leased":
leased.append(item)
elif item.traffic_state == "terminal_complete":
terminal_complete.append(item)
elif item.traffic_state == "blocked":
blocked.append(item)
elif item.traffic_state == "needs_controller":
needs_controller.append(item)
else:
runnable.append(item)
next_roles = [dashboard.next_safe_by_role[r].as_dict() for r in DASHBOARD_ROLES if r in dashboard.next_safe_by_role]
return TrafficSnapshot(
project_id=project_id,
repo_label=repo_label,
runnable=tuple(runnable),
leased=tuple(leased),
blocked=tuple(blocked),
needs_controller=tuple(needs_controller),
terminal_complete=tuple(terminal_complete),
next_roles=tuple(next_roles),
fetch_error=None,
inventory_complete=dashboard.inventory_complete,
)
def snapshot_to_dict(snapshot: TrafficSnapshot) -> dict[str, Any]:
return snapshot.as_dict()
+170
View File
@@ -0,0 +1,170 @@
"""HTML rendering for Phase 1 Traffic-Control View (#640)."""
from __future__ import annotations
from html import escape
from typing import Sequence
from webui.layout import render_page
from webui.traffic_loader import TrafficItem, TrafficSnapshot
def _render_badges(badges: Sequence[str]) -> str:
if not badges:
return ""
out = []
for b in badges:
cls = "badge"
b_lower = b.lower()
if "blocked" in b_lower or "unmet" in b_lower:
cls += " badge-blocked"
elif "claimed" in b_lower or "in-progress" in b_lower or "leased" in b_lower:
cls += " badge-claimed"
elif "review" in b_lower or "ready" in b_lower:
cls += " badge-in-review"
elif "duplicate" in b_lower:
cls += " badge-duplicate"
elif "stale" in b_lower:
cls += " badge-stale"
out.append(f'<span class="{cls}">{escape(b)}</span>')
return f'<div class="badges">{"".join(out)}</div>'
def _render_traffic_item_row(item: TrafficItem) -> str:
kind_label = escape(item.kind.upper())
num_str = f"#{item.number}"
title_str = escape(item.title)
role_str = escape(item.expected_role)
badges_html = _render_badges(item.badges)
reason_html = ""
if item.block_reason:
reason_html = f'<div class="muted" style="font-size:0.82rem; margin-top:0.2rem;"><strong>Blocker:</strong> {escape(item.block_reason)}</div>'
lease_html = ""
if item.lease_info:
owner = escape(str(item.lease_info.get("session_id") or item.lease_info.get("reviewer_identity") or "active worker"))
lease_html = f'<div class="muted" style="font-size:0.82rem; margin-top:0.2rem;"><strong>Lease:</strong> {owner}</div>'
return f"""<tr>
<td><code>{kind_label} {num_str}</code></td>
<td>
<div><strong>{title_str}</strong> {badges_html}</div>
{reason_html}
{lease_html}
</td>
<td><code>{role_str}</code></td>
</tr>"""
def _render_traffic_table(items: Sequence[TrafficItem], empty_message: str) -> str:
if not items:
return f'<p class="muted">{escape(empty_message)}</p>'
rows = "".join(_render_traffic_item_row(item) for item in items)
return f"""<table class="registry">
<thead>
<tr>
<th style="width: 15%;">Item</th>
<th style="width: 65%;">Title & Details</th>
<th style="width: 20%;">Next Role</th>
</tr>
</thead>
<tbody>
{rows}
</tbody>
</table>"""
def _render_next_roles(next_roles: Sequence[dict]) -> str:
if not next_roles:
return ""
cards = []
for r in next_roles:
role = escape(r.get("role", "unknown"))
status = r.get("status", "idle")
prompt = escape(r.get("prompt", ""))
status_cls = "badge-health-ok" if status == "safe" else ("badge-blocked" if "blocked" in status else "badge-health-skipped")
cards.append(f"""<div class="health-card" style="margin-bottom:0.75rem;">
<div style="display:flex; justify-content:space-between; align-items:center;">
<h3>Role: <code>{role}</code></h3>
<span class="badge {status_cls}">status: {escape(status)}</span>
</div>
<p class="meta" style="margin:0.35rem 0 0;">{prompt}</p>
</div>""")
return f"""<div style="margin: 1.5rem 0;">
<h3>Next Safe Role Actions</h3>
{"".join(cards)}
</div>"""
def render_traffic_page(snapshot: TrafficSnapshot) -> str:
"""Render the full HTML view for workflow traffic control."""
if snapshot.fetch_error:
body = f"""<h2>Workflow Traffic Control</h2>
<p class="meta">Repository: <code>{escape(snapshot.repo_label)}</code></p>
<div class="health-card health-stale">
<h3>Traffic data unavailable</h3>
<p class="health-headline">{escape(snapshot.fetch_error)}</p>
<p class="muted">Fail closed: traffic state cannot be established cleanly. Check credentials or remote connectivity.</p>
</div>"""
return render_page(title="Traffic Control", body_html=body)
runnable_count = len(snapshot.runnable)
leased_count = len(snapshot.leased)
blocked_count = len(snapshot.blocked)
controller_count = len(snapshot.needs_controller)
terminal_count = len(snapshot.terminal_complete)
summary_bar = f"""<div class="health-card" style="display:flex; flex-wrap:wrap; gap:1rem; align-items:center;">
<div><strong>Runnable:</strong> <span class="badge badge-health-ok">{runnable_count}</span></div>
<div><strong>Leased:</strong> <span class="badge badge-claimed">{leased_count}</span></div>
<div><strong>Blocked:</strong> <span class="badge badge-blocked">{blocked_count}</span></div>
<div><strong>Needs Controller:</strong> <span class="badge badge-duplicate">{controller_count}</span></div>
<div><strong>Terminal Complete:</strong> <span class="badge badge-stale">{terminal_count}</span></div>
</div>"""
next_roles_html = _render_next_roles(snapshot.next_roles)
sections_html = f"""
<div class="prompt-card">
<h3>1. Runnable Lanes (Ready for Allocation)</h3>
<p class="muted">Safe work items with no unmet dependencies or active leases. Safe for allocation.</p>
{_render_traffic_table(snapshot.runnable, "No runnable items ready for allocation.")}
</div>
<div class="prompt-card">
<h3>2. In-Progress Work (Active Leases)</h3>
<p class="muted">Work items currently leased and actively being worked by an assigned role session.</p>
{_render_traffic_table(snapshot.leased, "No active leases in flight.")}
</div>
<div class="prompt-card">
<h3>3. Blocked Items (Dependencies / Locks)</h3>
<p class="muted">Items blocked by unmet dependency issues, 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)