Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1948d3dc21 | ||
|
|
069a9af7e6 |
@@ -57,6 +57,33 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
|
||||
| `/system-health` | System-health dashboard — readiness, version/uptime, dependencies, MCP namespaces, stale-runtime parity (#639) |
|
||||
| `/queue` | Live PR and issue queue dashboard (#429) |
|
||||
| `/api/queue` | JSON queue export with pagination metadata |
|
||||
| `/traffic` | Workflow traffic-control view — runnable, leased, blocked, needs-controller, terminal-complete (#640) |
|
||||
| `/api/traffic` | JSON traffic-control export with state classifications and next safe role actions |
|
||||
|
||||
### Traffic-control state vocabulary (#640)
|
||||
|
||||
The traffic view classifies each open issue/PR into exactly one bucket:
|
||||
|
||||
| Bucket | Meaning | Operator implication |
|
||||
|--------|---------|----------------------|
|
||||
| **runnable** | No active lease, no block reason, safe for its expected role | Next role may start work |
|
||||
| **leased** | Active author claim or reviewer PR lease | Do not stomp; wait or adopt via role tools |
|
||||
| **blocked** | Dependency, missing head pin, conflict, or status:blocked | Author/controller remediation first |
|
||||
| **needs_controller** | Contaminated / controller-only diagnosis | Controller only |
|
||||
| **terminal_complete** | Reconciler / terminal-lock territory | Reconciler cleanup path |
|
||||
|
||||
**Live path contracts (do not invent):**
|
||||
|
||||
- PR head pins come from `QueueItem.signals["head_sha"]` (full SHA). Display
|
||||
`extra["head_sha"]` is truncated and must never be used for routing.
|
||||
- Reviewer leases are keyed as `(pr, pr_number)` only — never via a linked
|
||||
`issue_number` on the same lease marker.
|
||||
- Issue claims come from `claim_inventory["entries"]`
|
||||
(`issue_claim_heartbeat.build_claim_inventory`). There is no `active_claims`
|
||||
key.
|
||||
- Queue display badges are only: `blocked`, `claimed`, `duplicate`, `stale`,
|
||||
`in-review`, `open`. Review verdicts (`request-changes`, `approved`) are
|
||||
**not** queue badges; traffic does not invent them from the queue loader.
|
||||
| `/projects` | Project registry list with status and onboarding progress (#427, #635) |
|
||||
| `/projects/{id}` | Project detail + onboarding checklist |
|
||||
| `/api/v1/projects` | Versioned JSON registry export (#635) |
|
||||
|
||||
-726
@@ -1,726 +0,0 @@
|
||||
"""Pre-restart drain proof and hard gate (#661).
|
||||
|
||||
A sanctioned MCP restart may only proceed after a machine-verifiable *drain
|
||||
proof* attests that unsafe work is clear and checkpoints are complete. Drain
|
||||
mode alone (#659) is not enough: without a proof, an apply path could still
|
||||
restart on a stale or false "ready" claim, dropping mutations and orphaning
|
||||
leases. This module defines the :class:`DrainProof` artifact, a fail-closed
|
||||
verifier, and the hard gate the sanctioned restart-apply path must consult.
|
||||
|
||||
Design rules (mirror :mod:`restart_coordinator` / :mod:`lease_lifecycle`):
|
||||
|
||||
* **Pure classification.** Every function here operates on already-gathered
|
||||
inputs and returns a structured result. Nothing touches the network, the
|
||||
filesystem, or a live process, so multi-session fixtures drive every branch.
|
||||
This module never restarts anything; the gate only *authorizes or denies*.
|
||||
* **Fail closed.** A missing, expired, tampered, or unclean proof denies the
|
||||
restart. Unknown checkpoint completeness is treated as *not complete*. An
|
||||
incomplete impact report can never yield a clean proof.
|
||||
* **Non-forgeable within the process.** The proof id is a keyed hash over the
|
||||
canonical proof contents using a per-process secret. A worker session cannot
|
||||
hand-craft a passing proof without that secret, and a proof minted in a prior
|
||||
daemon process will not verify after a restart (the secret is regenerated).
|
||||
* **No secrets leak.** The per-process secret never appears in a proof, an
|
||||
``as_dict``, an audit record, or an incident descriptor.
|
||||
|
||||
Relationship to siblings (#655 umbrella):
|
||||
|
||||
* **#658** ``restart_coordinator.evaluate_restart_impact`` — produces the
|
||||
blast-radius impact report this proof consumes ("what would a restart
|
||||
disrupt?"). ``mutations`` / ``critical_sections`` being empty is what the
|
||||
no-in-flight-mutations check verifies.
|
||||
* **#659** graceful drain mode — performs the drain actions and calls
|
||||
:func:`build_drain_proof` to mint the artifact once its checklist passes.
|
||||
* **#660** durable session checkpoints — supplies checkpoint completeness.
|
||||
Because that schema may not yet be present, completeness is an *input* here,
|
||||
never a hard table dependency; unknown fails closed.
|
||||
|
||||
Non-goals (separate children): the emergency break-glass *workflow* (this gate
|
||||
only leaves a sanctioned bypass hole authorized elsewhere), the console UI, and
|
||||
the actual restart execution.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import hmac
|
||||
import json
|
||||
import os
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Any, Mapping, Sequence
|
||||
|
||||
DRAIN_PROOF_VERSION = "1.0.0-issue-661"
|
||||
|
||||
# Short default lifetime for a drain proof. A proof attests to a *point-in-time*
|
||||
# drained state; live work can resume the moment drain mode relaxes, so the
|
||||
# window in which a proof is honoured must be small (#661 security: short TTL).
|
||||
DEFAULT_PROOF_TTL_SECONDS = 120
|
||||
|
||||
# Gate verdicts.
|
||||
GATE_ALLOW = "allow"
|
||||
GATE_DENY = "deny"
|
||||
GATE_BREAK_GLASS = "break_glass"
|
||||
|
||||
# The mandatory drain checklist. A proof is *clean* only when every one of these
|
||||
# checks passed. Names are stable identifiers surfaced in audit + incidents.
|
||||
CHECK_NO_INFLIGHT_MUTATIONS = "no_inflight_mutations"
|
||||
CHECK_ASSIGNMENTS_STOPPED = "assignments_stopped"
|
||||
CHECK_CHECKPOINTS_COMPLETE = "checkpoints_complete"
|
||||
CHECK_HANDOFFS_OK = "handoffs_ok"
|
||||
CHECK_LEASES_HANDLED = "leases_handled"
|
||||
CHECK_ACKS_OR_TIMEOUT = "acks_or_timeout"
|
||||
|
||||
REQUIRED_CHECKS: tuple[str, ...] = (
|
||||
CHECK_NO_INFLIGHT_MUTATIONS,
|
||||
CHECK_ASSIGNMENTS_STOPPED,
|
||||
CHECK_CHECKPOINTS_COMPLETE,
|
||||
CHECK_HANDOFFS_OK,
|
||||
CHECK_LEASES_HANDLED,
|
||||
CHECK_ACKS_OR_TIMEOUT,
|
||||
)
|
||||
|
||||
|
||||
def _utc_now() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def _parse_ts(value: str | None) -> datetime | None:
|
||||
if not value:
|
||||
return None
|
||||
text = str(value).strip()
|
||||
if not text:
|
||||
return None
|
||||
if text.endswith("Z"):
|
||||
text = text[:-1] + "+00:00"
|
||||
try:
|
||||
dt = datetime.fromisoformat(text)
|
||||
except ValueError:
|
||||
return None
|
||||
if dt.tzinfo is None:
|
||||
dt = dt.replace(tzinfo=timezone.utc)
|
||||
return dt
|
||||
|
||||
|
||||
def _canonical(payload: Any) -> str:
|
||||
"""Deterministic JSON encoding for hashing (stable key order, no spaces)."""
|
||||
|
||||
return json.dumps(payload, sort_keys=True, separators=(",", ":"), default=str)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Per-process secret. Generated once per daemon process; regenerated on restart.
|
||||
# Injectable for tests so build + verify share a secret. Never serialized.
|
||||
# ---------------------------------------------------------------------------
|
||||
_PROCESS_SECRET = os.urandom(32)
|
||||
|
||||
|
||||
def process_secret() -> bytes:
|
||||
"""Return the per-process proof-signing secret (never serialized)."""
|
||||
|
||||
return _PROCESS_SECRET
|
||||
|
||||
|
||||
def _resolve_secret(secret: bytes | None) -> bytes:
|
||||
return secret if secret is not None else _PROCESS_SECRET
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class DrainCheck:
|
||||
"""One mandatory drain checklist result."""
|
||||
|
||||
name: str
|
||||
passed: bool
|
||||
detail: str
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {"name": self.name, "passed": self.passed, "detail": self.detail}
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class DrainProof:
|
||||
"""Machine-verifiable proof that a restart's blast radius has been drained.
|
||||
|
||||
The ``proof_id`` is a keyed hash over the canonical proof body; it is the
|
||||
tamper-evident signature verified at the gate. ``clean`` is True only when
|
||||
every required check passed. The proof is honoured only until ``expires_at``.
|
||||
"""
|
||||
|
||||
version: str
|
||||
proof_id: str
|
||||
clean: bool
|
||||
issued_at: str
|
||||
expires_at: str
|
||||
requesting_session_id: str | None
|
||||
impact_fingerprint: str
|
||||
checks: list[DrainCheck]
|
||||
failed_checks: list[str]
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
"version": self.version,
|
||||
"proof_id": self.proof_id,
|
||||
"clean": self.clean,
|
||||
"issued_at": self.issued_at,
|
||||
"expires_at": self.expires_at,
|
||||
"requesting_session_id": self.requesting_session_id,
|
||||
"impact_fingerprint": self.impact_fingerprint,
|
||||
"checks": [c.as_dict() for c in self.checks],
|
||||
"failed_checks": list(self.failed_checks),
|
||||
}
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class VerifyResult:
|
||||
"""Outcome of :func:`verify_drain_proof` (fail closed)."""
|
||||
|
||||
valid: bool
|
||||
reasons: list[str]
|
||||
proof_id: str | None
|
||||
clean: bool
|
||||
expired: bool
|
||||
tampered: bool
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
"valid": self.valid,
|
||||
"reasons": list(self.reasons),
|
||||
"proof_id": self.proof_id,
|
||||
"clean": self.clean,
|
||||
"expired": self.expired,
|
||||
"tampered": self.tampered,
|
||||
}
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class GateDecision:
|
||||
"""Outcome of :func:`gate_apply_restart`."""
|
||||
|
||||
allow: bool
|
||||
verdict: str
|
||||
reasons: list[str]
|
||||
proof_id: str | None
|
||||
break_glass: bool
|
||||
incident: dict[str, Any] | None
|
||||
audit_record: dict[str, Any]
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
"allow": self.allow,
|
||||
"verdict": self.verdict,
|
||||
"reasons": list(self.reasons),
|
||||
"proof_id": self.proof_id,
|
||||
"break_glass": self.break_glass,
|
||||
"incident": self.incident,
|
||||
"audit_record": dict(self.audit_record),
|
||||
}
|
||||
|
||||
|
||||
def impact_fingerprint(impact_report: Mapping[str, Any] | None) -> str:
|
||||
"""Stable fingerprint of the blast-radius state a proof was minted against.
|
||||
|
||||
Binds a proof to the specific impact evaluation. If the live state changes
|
||||
(a new mutation appears) between minting and gate, the caller can pass the
|
||||
fresh fingerprint and the proof will be rejected as stale.
|
||||
"""
|
||||
|
||||
report = impact_report or {}
|
||||
counts = report.get("counts") or {}
|
||||
material = {
|
||||
"inventory_complete": bool(report.get("inventory_complete", False)),
|
||||
"verdict": report.get("verdict"),
|
||||
"affected_issues": sorted(report.get("affected_issues") or []),
|
||||
"affected_prs": sorted(report.get("affected_prs") or []),
|
||||
"mutations": sorted(
|
||||
str(m.get("lease_id"))
|
||||
for m in (report.get("mutations") or [])
|
||||
if isinstance(m, Mapping)
|
||||
),
|
||||
"critical_sections": sorted(
|
||||
str(c.get("lease_id"))
|
||||
for c in (report.get("critical_sections") or [])
|
||||
if isinstance(c, Mapping)
|
||||
),
|
||||
"terminal_lock": bool(report.get("terminal_lock")),
|
||||
"counts": {
|
||||
k: counts.get(k)
|
||||
for k in (
|
||||
"sessions_live_other",
|
||||
"leases_disruptive",
|
||||
"critical_sections",
|
||||
"mutations",
|
||||
)
|
||||
},
|
||||
}
|
||||
return hashlib.sha256(_canonical(material).encode("utf-8")).hexdigest()
|
||||
|
||||
|
||||
def _sign(body: Mapping[str, Any], secret: bytes) -> str:
|
||||
"""Keyed (HMAC-SHA256) signature over the canonical proof body."""
|
||||
|
||||
return hmac.new(
|
||||
secret, _canonical(body).encode("utf-8"), hashlib.sha256
|
||||
).hexdigest()
|
||||
|
||||
|
||||
def _proof_body(
|
||||
*,
|
||||
clean: bool,
|
||||
issued_at: str,
|
||||
expires_at: str,
|
||||
requesting_session_id: str | None,
|
||||
fingerprint: str,
|
||||
checks: Sequence[DrainCheck],
|
||||
) -> dict[str, Any]:
|
||||
"""The exact fields covered by the signature. Order-independent (canonical)."""
|
||||
|
||||
return {
|
||||
"version": DRAIN_PROOF_VERSION,
|
||||
"clean": clean,
|
||||
"issued_at": issued_at,
|
||||
"expires_at": expires_at,
|
||||
"requesting_session_id": requesting_session_id,
|
||||
"impact_fingerprint": fingerprint,
|
||||
"checks": [c.as_dict() for c in checks],
|
||||
}
|
||||
|
||||
|
||||
def _bool_input(value: Any) -> bool:
|
||||
"""Strictly interpret a drain-state flag; anything not explicitly True fails."""
|
||||
|
||||
return value is True
|
||||
|
||||
|
||||
def _evaluate_checks(
|
||||
impact_report: Mapping[str, Any],
|
||||
drain_state: Mapping[str, Any],
|
||||
) -> list[DrainCheck]:
|
||||
"""Compute the mandatory checklist from the impact report + drain outcomes.
|
||||
|
||||
The report answers "is unsafe work still in flight?"; ``drain_state`` reports
|
||||
the drain-mode actions the coordinator/#659 performed. Every check fails
|
||||
closed when its evidence is missing or ambiguous.
|
||||
"""
|
||||
|
||||
checks: list[DrainCheck] = []
|
||||
|
||||
inventory_complete = bool(impact_report.get("inventory_complete", False))
|
||||
mutations = list(impact_report.get("mutations") or [])
|
||||
critical = list(impact_report.get("critical_sections") or [])
|
||||
disruptive = [
|
||||
l
|
||||
for l in (impact_report.get("affected_leases") or [])
|
||||
if isinstance(l, Mapping) and l.get("disruptive")
|
||||
]
|
||||
|
||||
# 1. No in-flight mutations. Derived from the authoritative impact report,
|
||||
# not self-reported: a proof cannot claim "no mutations" while the report
|
||||
# still shows mutations or unsevered critical sections.
|
||||
if not inventory_complete:
|
||||
checks.append(
|
||||
DrainCheck(
|
||||
CHECK_NO_INFLIGHT_MUTATIONS,
|
||||
False,
|
||||
"impact report inventory incomplete; cannot confirm mutations "
|
||||
"cleared (fail closed)",
|
||||
)
|
||||
)
|
||||
elif mutations or critical:
|
||||
checks.append(
|
||||
DrainCheck(
|
||||
CHECK_NO_INFLIGHT_MUTATIONS,
|
||||
False,
|
||||
f"{len(mutations)} mutation(s) and {len(critical)} critical "
|
||||
"section(s) still in flight",
|
||||
)
|
||||
)
|
||||
else:
|
||||
checks.append(
|
||||
DrainCheck(
|
||||
CHECK_NO_INFLIGHT_MUTATIONS,
|
||||
True,
|
||||
"no in-flight mutations or critical sections in impact report",
|
||||
)
|
||||
)
|
||||
|
||||
# 2. New assignment halted (maintenance-drain entered).
|
||||
checks.append(
|
||||
DrainCheck(
|
||||
CHECK_ASSIGNMENTS_STOPPED,
|
||||
_bool_input(drain_state.get("assignments_stopped")),
|
||||
"new work assignment halted"
|
||||
if _bool_input(drain_state.get("assignments_stopped"))
|
||||
else "assignments not confirmed stopped (fail closed)",
|
||||
)
|
||||
)
|
||||
|
||||
# 3. Durable checkpoints complete (#660). Unknown => not complete.
|
||||
cp = drain_state.get("checkpoints_complete")
|
||||
checks.append(
|
||||
DrainCheck(
|
||||
CHECK_CHECKPOINTS_COMPLETE,
|
||||
_bool_input(cp),
|
||||
"all live sessions checkpointed"
|
||||
if _bool_input(cp)
|
||||
else "checkpoint completeness unconfirmed (fail closed)",
|
||||
)
|
||||
)
|
||||
|
||||
# 4. Handoffs verified.
|
||||
checks.append(
|
||||
DrainCheck(
|
||||
CHECK_HANDOFFS_OK,
|
||||
_bool_input(drain_state.get("handoffs_verified")),
|
||||
"pending handoffs verified"
|
||||
if _bool_input(drain_state.get("handoffs_verified"))
|
||||
else "handoffs not verified (fail closed)",
|
||||
)
|
||||
)
|
||||
|
||||
# 5. Leases resolved/transferred/preserved AND none left disruptive. Requires
|
||||
# both the drain-mode assertion and the report showing no disruptive lease.
|
||||
leases_asserted = _bool_input(drain_state.get("leases_handled"))
|
||||
if not leases_asserted:
|
||||
checks.append(
|
||||
DrainCheck(
|
||||
CHECK_LEASES_HANDLED,
|
||||
False,
|
||||
"lease disposition not asserted by drain (fail closed)",
|
||||
)
|
||||
)
|
||||
elif disruptive:
|
||||
checks.append(
|
||||
DrainCheck(
|
||||
CHECK_LEASES_HANDLED,
|
||||
False,
|
||||
f"{len(disruptive)} disruptive lease(s) still active in report",
|
||||
)
|
||||
)
|
||||
else:
|
||||
checks.append(
|
||||
DrainCheck(
|
||||
CHECK_LEASES_HANDLED,
|
||||
True,
|
||||
"leases resolved/transferred/preserved; none left disruptive",
|
||||
)
|
||||
)
|
||||
|
||||
# 6. Acknowledgements received, or an explicit timeout policy was applied.
|
||||
acks = drain_state.get("acks") or {}
|
||||
ack_values = list(acks.values()) if isinstance(acks, Mapping) else []
|
||||
all_acked = bool(ack_values) and all(
|
||||
str(v).strip().lower() in {"ack", "acked", "acknowledged"}
|
||||
for v in ack_values
|
||||
)
|
||||
no_sessions_to_ack = isinstance(acks, Mapping) and len(ack_values) == 0
|
||||
timeout_policy = _bool_input(drain_state.get("ack_timeout_policy_applied"))
|
||||
acks_ok = all_acked or no_sessions_to_ack or timeout_policy
|
||||
if acks_ok:
|
||||
if timeout_policy and not all_acked:
|
||||
detail = "explicit ack timeout policy applied"
|
||||
elif no_sessions_to_ack:
|
||||
detail = "no other live sessions required to acknowledge"
|
||||
else:
|
||||
detail = "all affected sessions acknowledged"
|
||||
else:
|
||||
detail = "outstanding acks with no timeout policy (fail closed)"
|
||||
checks.append(DrainCheck(CHECK_ACKS_OR_TIMEOUT, acks_ok, detail))
|
||||
|
||||
return checks
|
||||
|
||||
|
||||
def build_drain_proof(
|
||||
*,
|
||||
impact_report: Mapping[str, Any],
|
||||
drain_state: Mapping[str, Any],
|
||||
requesting_session_id: str | None = None,
|
||||
now: datetime | None = None,
|
||||
ttl_seconds: int = DEFAULT_PROOF_TTL_SECONDS,
|
||||
secret: bytes | None = None,
|
||||
) -> DrainProof:
|
||||
"""Mint a drain proof from an impact report and the drain-mode outcomes.
|
||||
|
||||
A successful drain (every checklist item passes) yields a *clean* proof with
|
||||
a valid signature (AC#2). An unclean drain still yields a signed proof, but
|
||||
with ``clean=False`` and the failing checks named — the gate will deny it —
|
||||
so the artifact is auditable rather than silently absent.
|
||||
"""
|
||||
|
||||
moment = now or _utc_now()
|
||||
ttl = max(1, int(ttl_seconds))
|
||||
issued_at = moment.isoformat()
|
||||
expires_at = (moment + timedelta(seconds=ttl)).isoformat()
|
||||
fingerprint = impact_fingerprint(impact_report)
|
||||
|
||||
checks = _evaluate_checks(impact_report, drain_state)
|
||||
failed = [c.name for c in checks if not c.passed]
|
||||
clean = not failed
|
||||
|
||||
body = _proof_body(
|
||||
clean=clean,
|
||||
issued_at=issued_at,
|
||||
expires_at=expires_at,
|
||||
requesting_session_id=requesting_session_id,
|
||||
fingerprint=fingerprint,
|
||||
checks=checks,
|
||||
)
|
||||
proof_id = _sign(body, _resolve_secret(secret))
|
||||
|
||||
return DrainProof(
|
||||
version=DRAIN_PROOF_VERSION,
|
||||
proof_id=proof_id,
|
||||
clean=clean,
|
||||
issued_at=issued_at,
|
||||
expires_at=expires_at,
|
||||
requesting_session_id=requesting_session_id,
|
||||
impact_fingerprint=fingerprint,
|
||||
checks=checks,
|
||||
failed_checks=failed,
|
||||
)
|
||||
|
||||
|
||||
def verify_drain_proof(
|
||||
proof: Mapping[str, Any] | None,
|
||||
*,
|
||||
now: datetime | None = None,
|
||||
secret: bytes | None = None,
|
||||
expected_impact_fingerprint: str | None = None,
|
||||
) -> VerifyResult:
|
||||
"""Verify a drain proof, failing closed on any doubt.
|
||||
|
||||
A proof is valid only when: it is present and well-formed; its signature
|
||||
recomputes with the per-process secret (not forged/tampered, not minted in a
|
||||
prior process); it has not expired; every required check is present and
|
||||
passed; and — when ``expected_impact_fingerprint`` is supplied — it was
|
||||
minted against the current blast-radius state.
|
||||
"""
|
||||
|
||||
moment = now or _utc_now()
|
||||
reasons: list[str] = []
|
||||
|
||||
if not isinstance(proof, Mapping):
|
||||
return VerifyResult(
|
||||
valid=False,
|
||||
reasons=["drain proof missing or not an object (fail closed)"],
|
||||
proof_id=None,
|
||||
clean=False,
|
||||
expired=False,
|
||||
tampered=False,
|
||||
)
|
||||
|
||||
proof_id = proof.get("proof_id")
|
||||
presented_clean = bool(proof.get("clean", False))
|
||||
|
||||
# Rebuild the signed body from the presented fields and re-sign. Any mutation
|
||||
# of a covered field (including flipping ``clean`` to True) breaks the match.
|
||||
raw_checks = proof.get("checks")
|
||||
checks: list[DrainCheck] = []
|
||||
checks_wellformed = isinstance(raw_checks, Sequence) and not isinstance(
|
||||
raw_checks, (str, bytes)
|
||||
)
|
||||
if checks_wellformed:
|
||||
for c in raw_checks:
|
||||
if not isinstance(c, Mapping) or "name" not in c or "passed" not in c:
|
||||
checks_wellformed = False
|
||||
break
|
||||
checks.append(
|
||||
DrainCheck(
|
||||
name=str(c.get("name")),
|
||||
passed=bool(c.get("passed")),
|
||||
detail=str(c.get("detail") or ""),
|
||||
)
|
||||
)
|
||||
|
||||
tampered = False
|
||||
if not checks_wellformed:
|
||||
reasons.append("drain proof checks malformed (fail closed)")
|
||||
tampered = True
|
||||
else:
|
||||
body = _proof_body(
|
||||
clean=presented_clean,
|
||||
issued_at=str(proof.get("issued_at") or ""),
|
||||
expires_at=str(proof.get("expires_at") or ""),
|
||||
requesting_session_id=proof.get("requesting_session_id"),
|
||||
fingerprint=str(proof.get("impact_fingerprint") or ""),
|
||||
checks=checks,
|
||||
)
|
||||
expected_sig = _sign(body, _resolve_secret(secret))
|
||||
if not (
|
||||
isinstance(proof_id, str)
|
||||
and hmac.compare_digest(expected_sig, proof_id)
|
||||
):
|
||||
tampered = True
|
||||
reasons.append(
|
||||
"drain proof signature mismatch: forged, tampered, or minted "
|
||||
"in a prior process (fail closed)"
|
||||
)
|
||||
|
||||
expires = _parse_ts(proof.get("expires_at"))
|
||||
expired = expires is None or moment >= expires
|
||||
if expires is None:
|
||||
reasons.append("drain proof has no valid expiry (fail closed)")
|
||||
elif expired:
|
||||
reasons.append(f"drain proof expired at {proof.get('expires_at')}")
|
||||
|
||||
# Recompute cleanliness from the checks themselves — never trust the flag.
|
||||
recomputed_failed = [c.name for c in checks if not c.passed]
|
||||
present_names = {c.name for c in checks}
|
||||
missing = [name for name in REQUIRED_CHECKS if name not in present_names]
|
||||
recomputed_clean = checks_wellformed and not recomputed_failed and not missing
|
||||
if missing:
|
||||
reasons.append(f"drain proof missing required checks: {', '.join(missing)}")
|
||||
if checks_wellformed and recomputed_failed:
|
||||
reasons.append(
|
||||
f"drain checks failed: {', '.join(sorted(set(recomputed_failed)))}"
|
||||
)
|
||||
if presented_clean and not recomputed_clean:
|
||||
tampered = True
|
||||
reasons.append("proof claims clean but its checks do not support it")
|
||||
|
||||
if expected_impact_fingerprint is not None:
|
||||
if str(proof.get("impact_fingerprint") or "") != str(
|
||||
expected_impact_fingerprint
|
||||
):
|
||||
reasons.append(
|
||||
"drain proof was minted against a different blast-radius state "
|
||||
"(stale; fail closed)"
|
||||
)
|
||||
|
||||
valid = (not tampered) and (not expired) and recomputed_clean and not reasons
|
||||
return VerifyResult(
|
||||
valid=valid,
|
||||
reasons=reasons,
|
||||
proof_id=proof_id if isinstance(proof_id, str) else None,
|
||||
clean=recomputed_clean,
|
||||
expired=expired,
|
||||
tampered=tampered,
|
||||
)
|
||||
|
||||
|
||||
def _incident_descriptor(
|
||||
*,
|
||||
reasons: Sequence[str],
|
||||
requesting_session_id: str | None,
|
||||
proof_id: str | None,
|
||||
at: str,
|
||||
) -> dict[str, Any]:
|
||||
"""Durable incident descriptor for a denied restart (caller creates issue).
|
||||
|
||||
Kept as data (not a live Gitea call) so this module stays pure and testable;
|
||||
the MCP tool layer turns it into a durable issue via the incident bridge.
|
||||
"""
|
||||
|
||||
return {
|
||||
"kind": "restart_drain_gate_denied",
|
||||
"title": "Restart denied: drain proof failed the hard gate",
|
||||
"labels": ["mcp-health", "safety", "stale-runtime", "workflow-hardening"],
|
||||
"reasons": list(reasons),
|
||||
"requesting_session_id": requesting_session_id,
|
||||
"proof_id": proof_id,
|
||||
"at": at,
|
||||
"drain_proof_version": DRAIN_PROOF_VERSION,
|
||||
}
|
||||
|
||||
|
||||
def gate_apply_restart(
|
||||
*,
|
||||
proof: Mapping[str, Any] | None,
|
||||
now: datetime | None = None,
|
||||
secret: bytes | None = None,
|
||||
break_glass: bool = False,
|
||||
expected_impact_fingerprint: str | None = None,
|
||||
requesting_session_id: str | None = None,
|
||||
) -> GateDecision:
|
||||
"""Hard gate for a sanctioned restart apply (#661 AC#1/#3).
|
||||
|
||||
Restart apply is authorized only with a valid, unexpired, clean drain proof.
|
||||
A missing, expired, tampered, or unclean proof denies the restart and emits
|
||||
a durable incident descriptor. ``break_glass`` is the *only* sanctioned
|
||||
bypass — authorization for it is the caller's responsibility (the emergency
|
||||
workflow is a separate child); when set, the gate allows without a proof but
|
||||
records the bypass so it is never silent.
|
||||
"""
|
||||
|
||||
moment = now or _utc_now()
|
||||
at = moment.isoformat()
|
||||
|
||||
if break_glass:
|
||||
reasons = ["break-glass restart authorized; drain proof gate bypassed"]
|
||||
audit = {
|
||||
"event": "restart_gate_evaluated",
|
||||
"drain_proof_version": DRAIN_PROOF_VERSION,
|
||||
"evaluated_at": at,
|
||||
"verdict": GATE_BREAK_GLASS,
|
||||
"allow": True,
|
||||
"break_glass": True,
|
||||
"requesting_session_id": requesting_session_id,
|
||||
"proof_id": None,
|
||||
}
|
||||
return GateDecision(
|
||||
allow=True,
|
||||
verdict=GATE_BREAK_GLASS,
|
||||
reasons=reasons,
|
||||
proof_id=None,
|
||||
break_glass=True,
|
||||
incident=None,
|
||||
audit_record=audit,
|
||||
)
|
||||
|
||||
result = verify_drain_proof(
|
||||
proof,
|
||||
now=moment,
|
||||
secret=secret,
|
||||
expected_impact_fingerprint=expected_impact_fingerprint,
|
||||
)
|
||||
|
||||
if result.valid:
|
||||
reasons = ["valid unexpired clean drain proof present; restart authorized"]
|
||||
audit = {
|
||||
"event": "restart_gate_evaluated",
|
||||
"drain_proof_version": DRAIN_PROOF_VERSION,
|
||||
"evaluated_at": at,
|
||||
"verdict": GATE_ALLOW,
|
||||
"allow": True,
|
||||
"break_glass": False,
|
||||
"requesting_session_id": requesting_session_id,
|
||||
"proof_id": result.proof_id,
|
||||
}
|
||||
return GateDecision(
|
||||
allow=True,
|
||||
verdict=GATE_ALLOW,
|
||||
reasons=reasons,
|
||||
proof_id=result.proof_id,
|
||||
break_glass=False,
|
||||
incident=None,
|
||||
audit_record=audit,
|
||||
)
|
||||
|
||||
deny_reasons = ["restart denied: drain proof invalid (fail closed)"] + list(
|
||||
result.reasons
|
||||
)
|
||||
incident = _incident_descriptor(
|
||||
reasons=deny_reasons,
|
||||
requesting_session_id=requesting_session_id,
|
||||
proof_id=result.proof_id,
|
||||
at=at,
|
||||
)
|
||||
audit = {
|
||||
"event": "restart_gate_evaluated",
|
||||
"drain_proof_version": DRAIN_PROOF_VERSION,
|
||||
"evaluated_at": at,
|
||||
"verdict": GATE_DENY,
|
||||
"allow": False,
|
||||
"break_glass": False,
|
||||
"requesting_session_id": requesting_session_id,
|
||||
"proof_id": result.proof_id,
|
||||
"incident_kind": incident["kind"],
|
||||
}
|
||||
return GateDecision(
|
||||
allow=False,
|
||||
verdict=GATE_DENY,
|
||||
reasons=deny_reasons,
|
||||
proof_id=result.proof_id,
|
||||
break_glass=False,
|
||||
incident=incident,
|
||||
audit_record=audit,
|
||||
)
|
||||
+8
-58
@@ -2068,7 +2068,6 @@ 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)
|
||||
@@ -22343,8 +22342,6 @@ 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).
|
||||
|
||||
@@ -22354,15 +22351,10 @@ 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 **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``.
|
||||
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.
|
||||
|
||||
Operator override authority is read from the process environment
|
||||
(``GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION``), never self-asserted by
|
||||
@@ -22481,54 +22473,12 @@ 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:
|
||||
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
|
||||
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)"
|
||||
]
|
||||
return payload
|
||||
|
||||
|
||||
|
||||
@@ -1,383 +0,0 @@
|
||||
"""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)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,408 @@
|
||||
"""Tests for web UI workflow traffic-control view (#640)."""
|
||||
|
||||
import sys
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from unittest import mock
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
|
||||
|
||||
from starlette.testclient import TestClient
|
||||
|
||||
from webui.app import create_app
|
||||
from webui.traffic_loader import (
|
||||
TrafficItem,
|
||||
TrafficSnapshot,
|
||||
load_traffic_snapshot,
|
||||
snapshot_to_dict,
|
||||
)
|
||||
from webui.traffic_views import render_traffic_page
|
||||
from allocator_service import WorkCandidate
|
||||
|
||||
|
||||
class TestTrafficClassification(unittest.TestCase):
|
||||
def test_runnable_candidate_classification(self):
|
||||
cand = WorkCandidate(
|
||||
kind="issue",
|
||||
number=640,
|
||||
state="open",
|
||||
labels=("status:ready",),
|
||||
title="Web Console: Workflow traffic-control view (Phase 1)",
|
||||
priority=20,
|
||||
)
|
||||
snap = load_traffic_snapshot(candidates=[cand])
|
||||
self.assertEqual(len(snap.runnable), 1)
|
||||
self.assertEqual(snap.runnable[0].number, 640)
|
||||
self.assertTrue(snap.runnable[0].is_safe)
|
||||
self.assertEqual(snap.runnable[0].traffic_state, "runnable")
|
||||
|
||||
def test_blocked_dependency_candidate_classification(self):
|
||||
cand = WorkCandidate(
|
||||
kind="issue",
|
||||
number=643,
|
||||
state="open",
|
||||
labels=("status:ready",),
|
||||
title="Web Console: Requests & intent preview (Phase 2)",
|
||||
priority=20,
|
||||
dependency_unmet=True,
|
||||
dependency_reason="issue#643 depends on unresolved issue(s) #640; they are not closed",
|
||||
)
|
||||
snap = load_traffic_snapshot(candidates=[cand])
|
||||
self.assertEqual(len(snap.blocked), 1)
|
||||
self.assertEqual(snap.blocked[0].number, 643)
|
||||
self.assertFalse(snap.blocked[0].is_safe)
|
||||
self.assertEqual(snap.blocked[0].traffic_state, "blocked")
|
||||
self.assertIn("depends on unresolved issue(s) #640", snap.blocked[0].block_reason)
|
||||
|
||||
def test_leased_candidate_classification(self):
|
||||
cand = WorkCandidate(
|
||||
kind="issue",
|
||||
number=640,
|
||||
state="open",
|
||||
labels=("status:in-progress",),
|
||||
title="Web Console: Workflow traffic-control view (Phase 1)",
|
||||
priority=20,
|
||||
)
|
||||
lease = {
|
||||
"kind": "issue",
|
||||
"number": 640,
|
||||
"session_id": "prgs-author-12345",
|
||||
"role": "author",
|
||||
"status": "active",
|
||||
}
|
||||
snap = load_traffic_snapshot(candidates=[cand], leases=[lease])
|
||||
self.assertEqual(len(snap.leased), 1)
|
||||
self.assertEqual(snap.leased[0].number, 640)
|
||||
self.assertEqual(snap.leased[0].traffic_state, "leased")
|
||||
self.assertIsNotNone(snap.leased[0].lease_info)
|
||||
|
||||
def test_needs_controller_candidate_classification(self):
|
||||
cand = WorkCandidate(
|
||||
kind="issue",
|
||||
number=700,
|
||||
state="open",
|
||||
labels=("status:blocked",),
|
||||
title="Controller intervention needed",
|
||||
priority=10,
|
||||
blocked=True,
|
||||
)
|
||||
snap = load_traffic_snapshot(candidates=[cand])
|
||||
self.assertEqual(len(snap.needs_controller), 1)
|
||||
self.assertEqual(snap.needs_controller[0].number, 700)
|
||||
|
||||
|
||||
class TestTrafficLoader(unittest.TestCase):
|
||||
def test_snapshot_to_dict_export(self):
|
||||
cand = WorkCandidate(
|
||||
kind="issue",
|
||||
number=640,
|
||||
state="open",
|
||||
labels=("status:ready",),
|
||||
title="Traffic control test",
|
||||
priority=20,
|
||||
)
|
||||
snap = load_traffic_snapshot(candidates=[cand])
|
||||
data = snapshot_to_dict(snap)
|
||||
self.assertEqual(data["project_id"], "gitea-tools")
|
||||
self.assertEqual(len(data["runnable"]), 1)
|
||||
self.assertTrue(data["inventory_complete"])
|
||||
|
||||
def test_fail_closed_error_handling(self):
|
||||
with mock.patch("webui.traffic_loader.load_queue_snapshot", side_effect=RuntimeError("Gitea connection failed")):
|
||||
snap = load_traffic_snapshot()
|
||||
self.assertIsNotNone(snap.fetch_error)
|
||||
self.assertIn("Failed to load traffic state", snap.fetch_error)
|
||||
self.assertEqual(len(snap.runnable), 0)
|
||||
self.assertFalse(snap.inventory_complete)
|
||||
|
||||
|
||||
class TestTrafficLivePath(unittest.TestCase):
|
||||
"""Live path tests: inject QueueSnapshot + LeaseSnapshot (no candidates=).
|
||||
|
||||
Covers the production ``load_traffic_snapshot()`` branch that ``/traffic``
|
||||
and ``/api/traffic`` actually execute (#640 B1–B5).
|
||||
"""
|
||||
|
||||
FULL_SHA = "069a9af7e6aa2c2994e07199d1b0814819457017"
|
||||
|
||||
def _queue(
|
||||
self,
|
||||
*,
|
||||
prs=(),
|
||||
issues=(),
|
||||
):
|
||||
from webui.queue_loader import QueueSnapshot
|
||||
|
||||
return QueueSnapshot(
|
||||
project_id="gitea-tools",
|
||||
repo_label="Scaled-Tech-Consulting/Gitea-Tools",
|
||||
prs=tuple(prs),
|
||||
issues=tuple(issues),
|
||||
pr_pagination=None,
|
||||
issue_pagination=None,
|
||||
fetch_error=None,
|
||||
)
|
||||
|
||||
def _lease(
|
||||
self,
|
||||
*,
|
||||
claim_inventory=None,
|
||||
reviewer_leases=(),
|
||||
):
|
||||
from webui.lease_loader import LeaseSnapshot
|
||||
|
||||
return LeaseSnapshot(
|
||||
project_id="gitea-tools",
|
||||
repo_label="Scaled-Tech-Consulting/Gitea-Tools",
|
||||
issue_lock=None,
|
||||
claim_inventory=claim_inventory or {"entries": [], "counts": {}},
|
||||
reviewer_leases=tuple(reviewer_leases),
|
||||
duplicate_prs=(),
|
||||
duplicate_branches=(),
|
||||
collision_history=(),
|
||||
fetch_error=None,
|
||||
)
|
||||
|
||||
def test_live_pr_uses_full_head_sha_and_is_runnable(self):
|
||||
from webui.queue_loader import QueueItem
|
||||
|
||||
pr = QueueItem(
|
||||
number=885,
|
||||
title="traffic control",
|
||||
badges=("in-review",),
|
||||
extra={"head_sha": self.FULL_SHA[:12], "linked_issue": "640"},
|
||||
signals={
|
||||
"head_sha": self.FULL_SHA,
|
||||
"mergeable": True,
|
||||
"labels": (),
|
||||
"linked_issue": 640,
|
||||
},
|
||||
)
|
||||
q = self._queue(prs=[pr])
|
||||
l = self._lease()
|
||||
snap = load_traffic_snapshot(
|
||||
fetch_queue_snapshot=lambda: q,
|
||||
fetch_lease_snapshot=lambda: l,
|
||||
)
|
||||
self.assertIsNone(snap.fetch_error)
|
||||
self.assertEqual(len(snap.runnable), 1)
|
||||
item = snap.runnable[0]
|
||||
self.assertEqual(item.kind, "pr")
|
||||
self.assertEqual(item.number, 885)
|
||||
self.assertEqual(item.head_sha, self.FULL_SHA)
|
||||
self.assertNotEqual(item.head_sha, self.FULL_SHA[:12])
|
||||
self.assertIsNone(item.block_reason)
|
||||
self.assertEqual(len(snap.blocked), 0)
|
||||
|
||||
def test_live_pr_without_head_sha_is_blocked(self):
|
||||
from webui.queue_loader import QueueItem
|
||||
|
||||
pr = QueueItem(
|
||||
number=1,
|
||||
title="missing pin",
|
||||
badges=("open",),
|
||||
extra={"head_sha": ""},
|
||||
signals={"head_sha": "", "mergeable": True, "labels": ()},
|
||||
)
|
||||
snap = load_traffic_snapshot(
|
||||
fetch_queue_snapshot=lambda: self._queue(prs=[pr]),
|
||||
fetch_lease_snapshot=lambda: self._lease(),
|
||||
)
|
||||
self.assertEqual(len(snap.blocked) + len(snap.needs_controller), 1)
|
||||
item = (snap.blocked or snap.needs_controller)[0]
|
||||
self.assertIn("missing head_sha", (item.block_reason or "").lower())
|
||||
|
||||
def test_reviewer_lease_keys_by_pr_not_linked_issue(self):
|
||||
from webui.queue_loader import QueueItem
|
||||
|
||||
pr = QueueItem(
|
||||
number=885,
|
||||
title="leased pr",
|
||||
badges=("in-review",),
|
||||
extra={"head_sha": self.FULL_SHA[:12]},
|
||||
signals={"head_sha": self.FULL_SHA, "mergeable": True, "labels": ()},
|
||||
)
|
||||
issue = QueueItem(
|
||||
number=640,
|
||||
title="linked issue",
|
||||
badges=("open",),
|
||||
extra={},
|
||||
signals={"labels": ()},
|
||||
)
|
||||
# Marker-shaped record: has both pr_number and issue_number; must
|
||||
# attach to the PR only (B2).
|
||||
reviewer_lease = {
|
||||
"pr_number": 885,
|
||||
"issue_number": 640,
|
||||
"phase": "validating",
|
||||
"reviewer_identity": "sysadmin",
|
||||
"session_id": "review-sess-1",
|
||||
}
|
||||
snap = load_traffic_snapshot(
|
||||
fetch_queue_snapshot=lambda: self._queue(prs=[pr], issues=[issue]),
|
||||
fetch_lease_snapshot=lambda: self._lease(reviewer_leases=[reviewer_lease]),
|
||||
)
|
||||
leased_prs = [i for i in snap.leased if i.kind == "pr" and i.number == 885]
|
||||
self.assertEqual(len(leased_prs), 1)
|
||||
self.assertEqual(leased_prs[0].lease_info.get("pr_number"), 885)
|
||||
# Issue 640 must not inherit the reviewer lease just because issue_number
|
||||
# is present on the marker.
|
||||
for item in list(snap.leased) + list(snap.runnable) + list(snap.blocked):
|
||||
if item.kind == "issue" and item.number == 640:
|
||||
self.assertIsNone(
|
||||
item.lease_info,
|
||||
"reviewer lease must not attach to linked issue #640",
|
||||
)
|
||||
break
|
||||
else:
|
||||
self.fail("expected issue #640 in traffic snapshot")
|
||||
|
||||
def test_claim_inventory_entries_key_marks_issue_leased(self):
|
||||
from webui.queue_loader import QueueItem
|
||||
|
||||
issue = QueueItem(
|
||||
number=640,
|
||||
title="claimed issue",
|
||||
badges=("claimed",),
|
||||
extra={},
|
||||
signals={"labels": ("status:in-progress",)},
|
||||
)
|
||||
inventory = {
|
||||
"entries": [
|
||||
{
|
||||
"issue_number": 640,
|
||||
"status": "active",
|
||||
"latest_heartbeat": {"session_id": "author-sess-9"},
|
||||
"reasons": ["claim has structured heartbeat proof"],
|
||||
}
|
||||
],
|
||||
"counts": {"active": 1},
|
||||
"in_progress_total": 1,
|
||||
}
|
||||
snap = load_traffic_snapshot(
|
||||
fetch_queue_snapshot=lambda: self._queue(issues=[issue]),
|
||||
fetch_lease_snapshot=lambda: self._lease(claim_inventory=inventory),
|
||||
)
|
||||
leased_issues = [i for i in snap.leased if i.kind == "issue" and i.number == 640]
|
||||
self.assertEqual(len(leased_issues), 1)
|
||||
self.assertEqual(leased_issues[0].traffic_state, "leased")
|
||||
|
||||
def test_active_claims_key_is_ignored(self):
|
||||
"""B3 regression: fictional ``active_claims`` must not create lease_info."""
|
||||
from webui.queue_loader import QueueItem
|
||||
|
||||
issue = QueueItem(
|
||||
number=640,
|
||||
title="open issue",
|
||||
badges=("open",),
|
||||
extra={},
|
||||
signals={"labels": ()},
|
||||
)
|
||||
# Only the broken key — must NOT produce lease_info. Entries-less
|
||||
# inventory is empty (entries is the real claim_inventory key).
|
||||
inventory = {
|
||||
"active_claims": [
|
||||
{
|
||||
"kind": "issue",
|
||||
"number": 640,
|
||||
"issue_number": 640,
|
||||
"status": "active",
|
||||
},
|
||||
],
|
||||
"counts": {},
|
||||
}
|
||||
snap = load_traffic_snapshot(
|
||||
fetch_queue_snapshot=lambda: self._queue(issues=[issue]),
|
||||
fetch_lease_snapshot=lambda: self._lease(claim_inventory=inventory),
|
||||
)
|
||||
items = [
|
||||
i
|
||||
for i in (
|
||||
list(snap.runnable)
|
||||
+ list(snap.leased)
|
||||
+ list(snap.blocked)
|
||||
+ list(snap.needs_controller)
|
||||
)
|
||||
if i.kind == "issue" and i.number == 640
|
||||
]
|
||||
self.assertEqual(len(items), 1)
|
||||
self.assertIsNone(
|
||||
items[0].lease_info,
|
||||
"active_claims is not a real inventory key; entries-only",
|
||||
)
|
||||
|
||||
|
||||
class TestTrafficRoutesAndRendering(unittest.TestCase):
|
||||
def setUp(self):
|
||||
self.client = TestClient(create_app())
|
||||
|
||||
def test_traffic_html_page_rendering(self):
|
||||
cand1 = WorkCandidate(
|
||||
kind="issue",
|
||||
number=640,
|
||||
state="open",
|
||||
labels=("status:ready",),
|
||||
title="Traffic View Implementation",
|
||||
priority=20,
|
||||
)
|
||||
cand2 = WorkCandidate(
|
||||
kind="issue",
|
||||
number=643,
|
||||
state="open",
|
||||
labels=("status:ready",),
|
||||
title="Dependent Feature",
|
||||
priority=20,
|
||||
dependency_unmet=True,
|
||||
dependency_reason="issue#643 depends on unresolved issue(s) #640; they are not closed",
|
||||
)
|
||||
snap = load_traffic_snapshot(candidates=[cand1, cand2])
|
||||
with mock.patch("webui.app.load_traffic_snapshot", return_value=snap):
|
||||
response = self.client.get("/traffic")
|
||||
|
||||
self.assertEqual(response.status_code, 200)
|
||||
self.assertIn("Workflow Traffic Control", response.text)
|
||||
self.assertIn("1. Runnable Lanes", response.text)
|
||||
self.assertIn("3. Blocked Items", response.text)
|
||||
self.assertIn("Traffic View Implementation", response.text)
|
||||
self.assertIn("depends on unresolved issue(s) #640", response.text)
|
||||
|
||||
def test_api_traffic_json_route(self):
|
||||
cand = WorkCandidate(
|
||||
kind="issue",
|
||||
number=640,
|
||||
state="open",
|
||||
labels=("status:ready",),
|
||||
title="Traffic View API Test",
|
||||
priority=20,
|
||||
)
|
||||
snap = load_traffic_snapshot(candidates=[cand])
|
||||
with mock.patch("webui.app.load_traffic_snapshot", return_value=snap):
|
||||
response = self.client.get("/api/traffic")
|
||||
|
||||
self.assertEqual(response.status_code, 200)
|
||||
data = response.json()
|
||||
self.assertEqual(data["project_id"], "gitea-tools")
|
||||
self.assertEqual(len(data["runnable"]), 1)
|
||||
self.assertEqual(data["runnable"][0]["number"], 640)
|
||||
|
||||
def test_render_traffic_fail_closed_page(self):
|
||||
snap = TrafficSnapshot(
|
||||
project_id="gitea-tools",
|
||||
repo_label="Scaled-Tech-Consulting/Gitea-Tools",
|
||||
runnable=(),
|
||||
leased=(),
|
||||
blocked=(),
|
||||
needs_controller=(),
|
||||
terminal_complete=(),
|
||||
next_roles=(),
|
||||
fetch_error="Gitea credentials unavailable for gitea.prgs.cc",
|
||||
inventory_complete=False,
|
||||
)
|
||||
html = render_traffic_page(snap)
|
||||
self.assertIn("Traffic data unavailable", html)
|
||||
self.assertIn("Fail closed", html)
|
||||
self.assertNotIn("1. Runnable Lanes", html)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -42,6 +42,8 @@ from webui.lease_loader import load_lease_snapshot, snapshot_to_dict as lease_sn
|
||||
from webui.lease_views import render_leases_page
|
||||
from webui.queue_loader import load_queue_snapshot, snapshot_to_dict as queue_snapshot_to_dict
|
||||
from webui.queue_views import render_queue_page
|
||||
from webui.traffic_loader import load_traffic_snapshot, snapshot_to_dict as traffic_snapshot_to_dict
|
||||
from webui.traffic_views import render_traffic_page
|
||||
from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as worktree_snapshot_to_dict
|
||||
from webui.worktree_views import render_worktrees_page
|
||||
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
|
||||
@@ -200,6 +202,15 @@ async def api_queue(_request: Request) -> JSONResponse:
|
||||
return JSONResponse(queue_snapshot_to_dict(load_queue_snapshot()))
|
||||
|
||||
|
||||
async def traffic(_request: Request) -> HTMLResponse:
|
||||
snapshot = load_traffic_snapshot()
|
||||
return HTMLResponse(render_traffic_page(snapshot))
|
||||
|
||||
|
||||
async def api_traffic(_request: Request) -> JSONResponse:
|
||||
return JSONResponse(traffic_snapshot_to_dict(load_traffic_snapshot()))
|
||||
|
||||
|
||||
def _load_project_registry() -> tuple[ProjectRegistry | None, RegistryError | None]:
|
||||
"""Load the registry, converting validation failure into a fail-closed pair."""
|
||||
try:
|
||||
@@ -736,6 +747,8 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
|
||||
Route("/system-health", system_health, methods=["GET"]),
|
||||
Route("/queue", queue, methods=["GET"]),
|
||||
Route("/api/queue", api_queue, methods=["GET"]),
|
||||
Route("/traffic", traffic, methods=["GET"]),
|
||||
Route("/api/traffic", api_traffic, methods=["GET"]),
|
||||
Route("/projects", projects, methods=["GET"]),
|
||||
Route("/projects/{project_id}", project_detail, methods=["GET"]),
|
||||
Route("/api/projects", api_projects, methods=["GET"]),
|
||||
|
||||
@@ -182,10 +182,17 @@ def _extract_reviewer_leases(
|
||||
parsed = parse_reviewer_lease_comment(comment.get("body") or "")
|
||||
if not parsed:
|
||||
continue
|
||||
subject_pr = parsed.get("pr_number") or pr_number
|
||||
leases.append(
|
||||
{
|
||||
**parsed,
|
||||
"pr_number": parsed.get("pr_number") or pr_number,
|
||||
"pr_number": subject_pr,
|
||||
# The lease subject is the PR, never the linked issue: a
|
||||
# reviewer lease on PR #N must not be attributed to issue #N
|
||||
# or to the issue that PR closes (#640).
|
||||
"kind": "pr",
|
||||
"number": subject_pr,
|
||||
"role": "reviewer",
|
||||
"comment_id": comment.get("id"),
|
||||
"author": (comment.get("user") or {}).get("login"),
|
||||
"created_at": comment.get("created_at"),
|
||||
|
||||
@@ -41,6 +41,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
|
||||
NavItem("/system-health", "System health"),
|
||||
)),
|
||||
NavGroup("Traffic", (
|
||||
NavItem("/traffic", "Traffic control"),
|
||||
NavItem("/queue", "Queue"),
|
||||
NavItem("/leases", "Leases"),
|
||||
NavItem("/actions", "Actions"),
|
||||
|
||||
+34
-3
@@ -4,7 +4,7 @@ from __future__ import annotations
|
||||
|
||||
import os
|
||||
import re
|
||||
from dataclasses import dataclass
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Callable
|
||||
from urllib.parse import urlparse
|
||||
@@ -31,10 +31,20 @@ class PaginationMeta:
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class QueueItem:
|
||||
"""One queue row.
|
||||
|
||||
``extra`` holds *display* strings for the queue page (values are truncated
|
||||
or humanized for rendering). ``signals`` holds the *authoritative* typed
|
||||
values taken straight from the Gitea payload, for consumers that classify
|
||||
or pin state rather than render it (#640). Never derive identity or
|
||||
concurrency decisions from ``extra``.
|
||||
"""
|
||||
|
||||
number: int
|
||||
title: str
|
||||
badges: tuple[str, ...]
|
||||
extra: dict[str, str]
|
||||
signals: dict[str, Any] = field(default_factory=dict)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
@@ -134,21 +144,37 @@ def _format_pr_item(pr: dict, badges: tuple[str, ...]) -> QueueItem:
|
||||
"mergeable" if mergeable is True else "conflicted" if mergeable is False else "unknown"
|
||||
)
|
||||
linked = _extract_linked_issue(pr.get("title"), pr.get("body"))
|
||||
head_sha = str(head.get("sha") or "")
|
||||
labels = tuple(
|
||||
str(lb.get("name") or "") for lb in (pr.get("labels") or []) if lb.get("name")
|
||||
)
|
||||
return QueueItem(
|
||||
number=int(pr["number"]),
|
||||
title=str(pr.get("title") or ""),
|
||||
badges=badges,
|
||||
extra={
|
||||
"branch": f"{head.get('ref', '?')} → {base.get('ref', '?')}",
|
||||
"head_sha": str(head.get("sha") or "")[:12],
|
||||
# Display only — truncated. Pin against signals["head_sha"] instead.
|
||||
"head_sha": head_sha[:12],
|
||||
"mergeable": merge_label,
|
||||
"linked_issue": str(linked) if linked is not None else "",
|
||||
},
|
||||
signals={
|
||||
"head_sha": head_sha,
|
||||
"head_ref": str(head.get("ref") or ""),
|
||||
"base_ref": str(base.get("ref") or ""),
|
||||
"mergeable": mergeable if isinstance(mergeable, bool) else None,
|
||||
"labels": labels,
|
||||
"linked_issue": linked,
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
|
||||
labels = ", ".join(lb.get("name", "") for lb in issue.get("labels", []))
|
||||
label_names = tuple(
|
||||
str(lb.get("name") or "") for lb in (issue.get("labels") or []) if lb.get("name")
|
||||
)
|
||||
labels = ", ".join(label_names)
|
||||
assignee = (issue.get("assignee") or {}).get("login", "")
|
||||
return QueueItem(
|
||||
number=int(issue["number"]),
|
||||
@@ -159,6 +185,11 @@ def _format_issue_item(issue: dict, badges: tuple[str, ...]) -> QueueItem:
|
||||
"assignee": assignee or "unassigned",
|
||||
"state": str(issue.get("state") or ""),
|
||||
},
|
||||
signals={
|
||||
"labels": label_names,
|
||||
"assignee": assignee,
|
||||
"state": str(issue.get("state") or ""),
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,449 @@
|
||||
"""Traffic-control view loader for Phase 1 operator web console (#640).
|
||||
|
||||
Combines queue snapshots, inventory leases, dependency graph classifications,
|
||||
and workflow dashboard rules to deliver full traffic-control visibility:
|
||||
runnable, leased (in-progress), blocked (dependency/lock), needs-controller,
|
||||
and terminal-complete candidates.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, Callable, Sequence
|
||||
|
||||
from webui.project_registry import find_project, load_registry
|
||||
from webui.queue_loader import load_queue_snapshot, QueueSnapshot
|
||||
from webui.lease_loader import load_lease_snapshot, LeaseSnapshot
|
||||
from workflow_dashboard import (
|
||||
DashboardSnapshot,
|
||||
QueueEntry,
|
||||
RoleNextAction,
|
||||
build_workflow_dashboard,
|
||||
DASHBOARD_ROLES,
|
||||
)
|
||||
from allocator_service import WorkCandidate
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class TrafficItem:
|
||||
kind: str # "issue" or "pr"
|
||||
number: int
|
||||
title: str
|
||||
traffic_state: str # "runnable", "leased", "blocked", "needs_controller", "terminal_complete"
|
||||
expected_role: str
|
||||
safe_for_roles: tuple[str, ...]
|
||||
badges: tuple[str, ...]
|
||||
block_reason: str | None = None
|
||||
lease_info: dict[str, Any] | None = None
|
||||
head_sha: str | None = None
|
||||
|
||||
@property
|
||||
def is_safe(self) -> bool:
|
||||
return self.block_reason is None and bool(self.safe_for_roles)
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
"kind": self.kind,
|
||||
"number": self.number,
|
||||
"title": self.title,
|
||||
"traffic_state": self.traffic_state,
|
||||
"expected_role": self.expected_role,
|
||||
"safe_for_roles": list(self.safe_for_roles),
|
||||
"badges": list(self.badges),
|
||||
"block_reason": self.block_reason,
|
||||
"lease_info": self.lease_info,
|
||||
"head_sha": self.head_sha,
|
||||
"is_safe": self.is_safe,
|
||||
}
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class TrafficSnapshot:
|
||||
project_id: str
|
||||
repo_label: str
|
||||
runnable: tuple[TrafficItem, ...]
|
||||
leased: tuple[TrafficItem, ...]
|
||||
blocked: tuple[TrafficItem, ...]
|
||||
needs_controller: tuple[TrafficItem, ...]
|
||||
terminal_complete: tuple[TrafficItem, ...]
|
||||
next_roles: tuple[dict[str, Any], ...]
|
||||
fetch_error: str | None = None
|
||||
inventory_complete: bool = True
|
||||
|
||||
def as_dict(self) -> dict[str, Any]:
|
||||
return {
|
||||
"project_id": self.project_id,
|
||||
"repo_label": self.repo_label,
|
||||
"runnable": [i.as_dict() for i in self.runnable],
|
||||
"leased": [i.as_dict() for i in self.leased],
|
||||
"blocked": [i.as_dict() for i in self.blocked],
|
||||
"needs_controller": [i.as_dict() for i in self.needs_controller],
|
||||
"terminal_complete": [i.as_dict() for i in self.terminal_complete],
|
||||
"next_roles": list(self.next_roles),
|
||||
"fetch_error": self.fetch_error,
|
||||
"inventory_complete": self.inventory_complete,
|
||||
}
|
||||
|
||||
|
||||
def _classify_traffic_item(
|
||||
entry: QueueEntry,
|
||||
*,
|
||||
lease_info: dict[str, Any] | None = None,
|
||||
) -> TrafficItem:
|
||||
"""Classify a QueueEntry into a TrafficItem with explicit traffic state."""
|
||||
badges = list(entry.badges)
|
||||
block_reason = entry.block_reason
|
||||
expected_role = entry.expected_role
|
||||
|
||||
entry_is_safe = entry.block_reason is None and bool(entry.safe_for_roles)
|
||||
# Lease state is checked first: an item that is both leased and blocked is
|
||||
# reported as leased. That is safe by construction — a leased item is never
|
||||
# placed in the runnable lane — and it keeps the operator's attention on the
|
||||
# session that currently owns the work. The blocker text still renders.
|
||||
if lease_info is not None or "in-progress" in badges or "claimed" in badges:
|
||||
state = "leased"
|
||||
elif expected_role == "reconciler" or "terminal-lock" in badges:
|
||||
state = "terminal_complete"
|
||||
elif expected_role == "controller" or "contaminated" in badges or "needs-controller" in badges:
|
||||
state = "needs_controller"
|
||||
elif (
|
||||
block_reason is not None
|
||||
or "blocked" in badges
|
||||
or "dependency-unmet" in badges
|
||||
or "blocked-by-terminal" in badges
|
||||
or "status:blocked" in badges
|
||||
):
|
||||
state = "blocked"
|
||||
elif entry_is_safe:
|
||||
state = "runnable"
|
||||
else:
|
||||
state = "needs_controller"
|
||||
|
||||
return TrafficItem(
|
||||
kind=entry.kind,
|
||||
number=entry.number,
|
||||
title=entry.title,
|
||||
traffic_state=state,
|
||||
expected_role=expected_role,
|
||||
safe_for_roles=entry.safe_for_roles,
|
||||
badges=tuple(badges),
|
||||
block_reason=block_reason,
|
||||
lease_info=lease_info,
|
||||
head_sha=entry.head_sha,
|
||||
)
|
||||
|
||||
|
||||
# Claim statuses from ``issue_claim_heartbeat.build_claim_inventory`` that mean
|
||||
# a live worker currently holds the issue. Everything else (``stale``,
|
||||
# ``phantom``, ``reclaimable``, ``not_claimed``) is reported through the
|
||||
# dashboard's stale-lease channel and is never rendered as an active lease.
|
||||
_ACTIVE_CLAIM_STATUSES = frozenset({"active", "awaiting_review"})
|
||||
|
||||
# Statuses that positively mean "not an active lease" for any lease record.
|
||||
_INACTIVE_LEASE_STATUSES = frozenset(
|
||||
{"expired", "stale", "released", "moot", "reclaimable", "phantom", "not_claimed"}
|
||||
)
|
||||
|
||||
|
||||
def _candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate]:
|
||||
"""Build allocator candidates from the queue loader's authoritative signals.
|
||||
|
||||
Display badges (``blocked``/``claimed``/``duplicate``/``stale``/
|
||||
``in-review``/``open``) are rendering hints, not routing state, so nothing
|
||||
here branches on them. Every routing field comes from
|
||||
``QueueItem.signals`` — the raw Gitea payload values.
|
||||
|
||||
The queue loader reads ``/pulls`` and ``/issues`` only; it never fetches
|
||||
review verdicts. ``request_changes_current_head`` / ``approval_on_current_head``
|
||||
are therefore left at their fail-safe ``False`` rather than being guessed
|
||||
from badges: an unproven approval must never route a PR to the merger.
|
||||
"""
|
||||
candidates: list[WorkCandidate] = []
|
||||
|
||||
for pr in q_snap.prs:
|
||||
signals = pr.signals or {}
|
||||
head_sha = str(signals.get("head_sha") or "").strip()
|
||||
mergeable = signals.get("mergeable")
|
||||
labels = tuple(str(x) for x in (signals.get("labels") or ()))
|
||||
candidates.append(
|
||||
WorkCandidate(
|
||||
kind="pr",
|
||||
number=pr.number,
|
||||
state="open",
|
||||
labels=labels,
|
||||
title=pr.title,
|
||||
# Full 40-char SHA from head.sha — never the 12-char display value.
|
||||
head_sha=head_sha or None,
|
||||
priority=5,
|
||||
mergeable=mergeable is True,
|
||||
blocked=mergeable is False or "status:blocked" in labels,
|
||||
)
|
||||
)
|
||||
|
||||
for issue in q_snap.issues:
|
||||
signals = issue.signals or {}
|
||||
labels = tuple(str(x) for x in (signals.get("labels") or ()))
|
||||
lowered = {label.lower() for label in labels}
|
||||
candidates.append(
|
||||
WorkCandidate(
|
||||
kind="issue",
|
||||
number=issue.number,
|
||||
state="open",
|
||||
labels=labels,
|
||||
title=issue.title,
|
||||
priority=20 if "status:ready" in lowered else 10,
|
||||
blocked="status:blocked" in lowered,
|
||||
# A live claim by another session is not this session's work.
|
||||
already_claimed_elsewhere="status:in-progress" in lowered,
|
||||
)
|
||||
)
|
||||
|
||||
return candidates
|
||||
|
||||
|
||||
def _claim_lease_records(inventory: dict[str, Any] | None) -> list[dict[str, Any]]:
|
||||
"""Normalize ``build_claim_inventory`` entries into lease records.
|
||||
|
||||
The inventory contract is ``{"entries", "counts", "heartbeat_lease_minutes",
|
||||
"reclaim_after_minutes", "in_progress_total"}``. Each entry is keyed by
|
||||
``issue_number``; the subject kind is therefore always ``issue``.
|
||||
"""
|
||||
entries = (inventory or {}).get("entries") or ()
|
||||
records: list[dict[str, Any]] = []
|
||||
for entry in entries:
|
||||
if not isinstance(entry, dict):
|
||||
continue
|
||||
number = entry.get("issue_number")
|
||||
if number is None:
|
||||
continue
|
||||
try:
|
||||
number_int = int(number)
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
heartbeat = entry.get("latest_heartbeat") or {}
|
||||
record = dict(entry)
|
||||
record.update(
|
||||
{
|
||||
"kind": "issue",
|
||||
"number": number_int,
|
||||
"role": "author",
|
||||
"lease_source": "issue-claim-heartbeat",
|
||||
}
|
||||
)
|
||||
if isinstance(heartbeat, dict):
|
||||
if heartbeat.get("session_id") and not record.get("session_id"):
|
||||
record["session_id"] = heartbeat.get("session_id")
|
||||
if heartbeat.get("author") and not record.get("author"):
|
||||
record["author"] = heartbeat.get("author")
|
||||
records.append(record)
|
||||
return records
|
||||
|
||||
|
||||
def _lease_subject(lease: dict[str, Any]) -> tuple[str, int] | None:
|
||||
"""Return the ``(kind, number)`` a lease record actually covers.
|
||||
|
||||
Fails closed: a record that does not identify exactly one subject is
|
||||
dropped rather than attributed to a guessed work item (#640 — never invent
|
||||
a lease, and never attach a PR lease to a same-numbered issue).
|
||||
"""
|
||||
kind = str(lease.get("kind") or lease.get("work_kind") or "").strip().lower()
|
||||
pr_number = lease.get("pr_number")
|
||||
issue_number = lease.get("issue_number")
|
||||
|
||||
if kind not in ("pr", "issue"):
|
||||
if pr_number is not None and issue_number is None:
|
||||
kind = "pr"
|
||||
elif issue_number is not None and pr_number is None:
|
||||
kind = "issue"
|
||||
else:
|
||||
return None
|
||||
|
||||
number = lease.get("number")
|
||||
if number is None:
|
||||
number = lease.get("work_number")
|
||||
if number is None:
|
||||
number = pr_number if kind == "pr" else issue_number
|
||||
if number is None:
|
||||
return None
|
||||
try:
|
||||
return kind, int(number)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _is_active_lease(lease: dict[str, Any]) -> bool:
|
||||
"""True when the record proves a worker currently holds the item."""
|
||||
if lease.get("stale") or lease.get("expired"):
|
||||
return False
|
||||
status = str(lease.get("status") or lease.get("lease_status") or "").strip().lower()
|
||||
if status in _INACTIVE_LEASE_STATUSES:
|
||||
return False
|
||||
if lease.get("lease_source") == "issue-claim-heartbeat":
|
||||
return status in _ACTIVE_CLAIM_STATUSES
|
||||
return True
|
||||
|
||||
|
||||
def load_traffic_snapshot(
|
||||
*,
|
||||
candidates: Sequence[WorkCandidate] | None = None,
|
||||
leases: Sequence[dict[str, Any]] | None = None,
|
||||
terminal_pr: int | None = None,
|
||||
fetch_queue_snapshot: Callable[[], QueueSnapshot] | None = None,
|
||||
fetch_lease_snapshot: Callable[[], LeaseSnapshot] | None = None,
|
||||
project_id: str = "gitea-tools",
|
||||
) -> TrafficSnapshot:
|
||||
"""Load and compute the traffic-control snapshot."""
|
||||
try:
|
||||
reg = load_registry()
|
||||
proj = find_project(reg, project_id)
|
||||
repo_label = proj.remote_repo if proj else "Scaled-Tech-Consulting/Gitea-Tools"
|
||||
except Exception:
|
||||
repo_label = "Scaled-Tech-Consulting/Gitea-Tools"
|
||||
|
||||
# Injected candidates path (pure unit testing)
|
||||
if candidates is not None:
|
||||
dashboard = build_workflow_dashboard(
|
||||
candidates=candidates,
|
||||
leases=leases,
|
||||
terminal_pr=terminal_pr,
|
||||
inventory_complete=True,
|
||||
)
|
||||
return _build_traffic_snapshot_from_dashboard(
|
||||
project_id=project_id,
|
||||
repo_label=repo_label,
|
||||
dashboard=dashboard,
|
||||
leases=leases or (),
|
||||
)
|
||||
|
||||
# Live snapshot loading
|
||||
q_loader = fetch_queue_snapshot or load_queue_snapshot
|
||||
l_loader = fetch_lease_snapshot or load_lease_snapshot
|
||||
|
||||
try:
|
||||
q_snap = q_loader()
|
||||
l_snap = l_loader()
|
||||
except Exception as exc: # noqa: BLE001
|
||||
return TrafficSnapshot(
|
||||
project_id=project_id,
|
||||
repo_label=repo_label,
|
||||
runnable=(),
|
||||
leased=(),
|
||||
blocked=(),
|
||||
needs_controller=(),
|
||||
terminal_complete=(),
|
||||
next_roles=(),
|
||||
fetch_error=f"Failed to load traffic state: {exc}",
|
||||
inventory_complete=False,
|
||||
)
|
||||
|
||||
if q_snap.fetch_error or l_snap.fetch_error:
|
||||
err = q_snap.fetch_error or l_snap.fetch_error
|
||||
return TrafficSnapshot(
|
||||
project_id=project_id,
|
||||
repo_label=repo_label,
|
||||
runnable=(),
|
||||
leased=(),
|
||||
blocked=(),
|
||||
needs_controller=(),
|
||||
terminal_complete=(),
|
||||
next_roles=(),
|
||||
fetch_error=err,
|
||||
inventory_complete=False,
|
||||
)
|
||||
|
||||
candidate_list = _candidates_from_queue_snapshot(q_snap)
|
||||
|
||||
raw_leases: list[dict[str, Any]] = _claim_lease_records(l_snap.claim_inventory)
|
||||
for r_lease in l_snap.reviewer_leases or ():
|
||||
if not isinstance(r_lease, dict):
|
||||
continue
|
||||
# Always pin reviewer leases to the PR subject, even if a linked
|
||||
# issue_number is present on the marker (#640 B2).
|
||||
normalized = dict(r_lease)
|
||||
subject = normalized.get("pr_number") or normalized.get("number")
|
||||
if subject is None:
|
||||
continue
|
||||
try:
|
||||
pr_num = int(subject)
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
normalized["kind"] = "pr"
|
||||
normalized["number"] = pr_num
|
||||
normalized["pr_number"] = pr_num
|
||||
normalized.setdefault("role", "reviewer")
|
||||
raw_leases.append(normalized)
|
||||
|
||||
dashboard = build_workflow_dashboard(
|
||||
candidates=candidate_list,
|
||||
leases=raw_leases,
|
||||
inventory_complete=q_snap.pr_pagination.inventory_complete if q_snap.pr_pagination else True,
|
||||
)
|
||||
|
||||
return _build_traffic_snapshot_from_dashboard(
|
||||
project_id=project_id,
|
||||
repo_label=repo_label,
|
||||
dashboard=dashboard,
|
||||
leases=raw_leases,
|
||||
)
|
||||
|
||||
|
||||
def _build_traffic_snapshot_from_dashboard(
|
||||
*,
|
||||
project_id: str,
|
||||
repo_label: str,
|
||||
dashboard: DashboardSnapshot,
|
||||
leases: Sequence[dict[str, Any]],
|
||||
) -> TrafficSnapshot:
|
||||
"""Classify dashboard entries into the 5 traffic state buckets."""
|
||||
all_entries = dashboard.open_prs + dashboard.open_issues
|
||||
|
||||
# Map each active lease onto the exact work item it covers. Records whose
|
||||
# subject cannot be determined, and claims that are stale/phantom/
|
||||
# reclaimable, are deliberately dropped instead of guessed.
|
||||
lease_map: dict[tuple[str, int], dict[str, Any]] = {}
|
||||
for lease in leases:
|
||||
if not isinstance(lease, dict) or not _is_active_lease(lease):
|
||||
continue
|
||||
subject = _lease_subject(lease)
|
||||
if subject is not None:
|
||||
lease_map[subject] = lease
|
||||
|
||||
runnable: list[TrafficItem] = []
|
||||
leased: list[TrafficItem] = []
|
||||
blocked: list[TrafficItem] = []
|
||||
needs_controller: list[TrafficItem] = []
|
||||
terminal_complete: list[TrafficItem] = []
|
||||
|
||||
for entry in all_entries:
|
||||
l_info = lease_map.get((entry.kind, entry.number))
|
||||
item = _classify_traffic_item(entry, lease_info=l_info)
|
||||
|
||||
if item.traffic_state == "leased":
|
||||
leased.append(item)
|
||||
elif item.traffic_state == "terminal_complete":
|
||||
terminal_complete.append(item)
|
||||
elif item.traffic_state == "blocked":
|
||||
blocked.append(item)
|
||||
elif item.traffic_state == "needs_controller":
|
||||
needs_controller.append(item)
|
||||
else:
|
||||
runnable.append(item)
|
||||
|
||||
next_roles = [dashboard.next_safe_by_role[r].as_dict() for r in DASHBOARD_ROLES if r in dashboard.next_safe_by_role]
|
||||
|
||||
return TrafficSnapshot(
|
||||
project_id=project_id,
|
||||
repo_label=repo_label,
|
||||
runnable=tuple(runnable),
|
||||
leased=tuple(leased),
|
||||
blocked=tuple(blocked),
|
||||
needs_controller=tuple(needs_controller),
|
||||
terminal_complete=tuple(terminal_complete),
|
||||
next_roles=tuple(next_roles),
|
||||
fetch_error=None,
|
||||
inventory_complete=dashboard.inventory_complete,
|
||||
)
|
||||
|
||||
|
||||
def snapshot_to_dict(snapshot: TrafficSnapshot) -> dict[str, Any]:
|
||||
return snapshot.as_dict()
|
||||
@@ -0,0 +1,170 @@
|
||||
"""HTML rendering for Phase 1 Traffic-Control View (#640)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from html import escape
|
||||
from typing import Sequence
|
||||
|
||||
from webui.layout import render_page
|
||||
from webui.traffic_loader import TrafficItem, TrafficSnapshot
|
||||
|
||||
|
||||
def _render_badges(badges: Sequence[str]) -> str:
|
||||
if not badges:
|
||||
return ""
|
||||
out = []
|
||||
for b in badges:
|
||||
cls = "badge"
|
||||
b_lower = b.lower()
|
||||
if "blocked" in b_lower or "unmet" in b_lower:
|
||||
cls += " badge-blocked"
|
||||
elif "claimed" in b_lower or "in-progress" in b_lower or "leased" in b_lower:
|
||||
cls += " badge-claimed"
|
||||
elif "review" in b_lower or "ready" in b_lower:
|
||||
cls += " badge-in-review"
|
||||
elif "duplicate" in b_lower:
|
||||
cls += " badge-duplicate"
|
||||
elif "stale" in b_lower:
|
||||
cls += " badge-stale"
|
||||
out.append(f'<span class="{cls}">{escape(b)}</span>')
|
||||
return f'<div class="badges">{"".join(out)}</div>'
|
||||
|
||||
|
||||
def _render_traffic_item_row(item: TrafficItem) -> str:
|
||||
kind_label = escape(item.kind.upper())
|
||||
num_str = f"#{item.number}"
|
||||
title_str = escape(item.title)
|
||||
role_str = escape(item.expected_role)
|
||||
badges_html = _render_badges(item.badges)
|
||||
|
||||
reason_html = ""
|
||||
if item.block_reason:
|
||||
reason_html = f'<div class="muted" style="font-size:0.82rem; margin-top:0.2rem;"><strong>Blocker:</strong> {escape(item.block_reason)}</div>'
|
||||
|
||||
lease_html = ""
|
||||
if item.lease_info:
|
||||
owner = escape(str(item.lease_info.get("session_id") or item.lease_info.get("reviewer_identity") or "active worker"))
|
||||
lease_html = f'<div class="muted" style="font-size:0.82rem; margin-top:0.2rem;"><strong>Lease:</strong> {owner}</div>'
|
||||
|
||||
return f"""<tr>
|
||||
<td><code>{kind_label} {num_str}</code></td>
|
||||
<td>
|
||||
<div><strong>{title_str}</strong> {badges_html}</div>
|
||||
{reason_html}
|
||||
{lease_html}
|
||||
</td>
|
||||
<td><code>{role_str}</code></td>
|
||||
</tr>"""
|
||||
|
||||
|
||||
def _render_traffic_table(items: Sequence[TrafficItem], empty_message: str) -> str:
|
||||
if not items:
|
||||
return f'<p class="muted">{escape(empty_message)}</p>'
|
||||
|
||||
rows = "".join(_render_traffic_item_row(item) for item in items)
|
||||
return f"""<table class="registry">
|
||||
<thead>
|
||||
<tr>
|
||||
<th style="width: 15%;">Item</th>
|
||||
<th style="width: 65%;">Title & Details</th>
|
||||
<th style="width: 20%;">Next Role</th>
|
||||
</tr>
|
||||
</thead>
|
||||
<tbody>
|
||||
{rows}
|
||||
</tbody>
|
||||
</table>"""
|
||||
|
||||
|
||||
def _render_next_roles(next_roles: Sequence[dict]) -> str:
|
||||
if not next_roles:
|
||||
return ""
|
||||
|
||||
cards = []
|
||||
for r in next_roles:
|
||||
role = escape(r.get("role", "unknown"))
|
||||
status = r.get("status", "idle")
|
||||
prompt = escape(r.get("prompt", ""))
|
||||
|
||||
status_cls = "badge-health-ok" if status == "safe" else ("badge-blocked" if "blocked" in status else "badge-health-skipped")
|
||||
cards.append(f"""<div class="health-card" style="margin-bottom:0.75rem;">
|
||||
<div style="display:flex; justify-content:space-between; align-items:center;">
|
||||
<h3>Role: <code>{role}</code></h3>
|
||||
<span class="badge {status_cls}">status: {escape(status)}</span>
|
||||
</div>
|
||||
<p class="meta" style="margin:0.35rem 0 0;">{prompt}</p>
|
||||
</div>""")
|
||||
|
||||
return f"""<div style="margin: 1.5rem 0;">
|
||||
<h3>Next Safe Role Actions</h3>
|
||||
{"".join(cards)}
|
||||
</div>"""
|
||||
|
||||
|
||||
def render_traffic_page(snapshot: TrafficSnapshot) -> str:
|
||||
"""Render the full HTML view for workflow traffic control."""
|
||||
if snapshot.fetch_error:
|
||||
body = f"""<h2>Workflow Traffic Control</h2>
|
||||
<p class="meta">Repository: <code>{escape(snapshot.repo_label)}</code></p>
|
||||
<div class="health-card health-stale">
|
||||
<h3>Traffic data unavailable</h3>
|
||||
<p class="health-headline">{escape(snapshot.fetch_error)}</p>
|
||||
<p class="muted">Fail closed: traffic state cannot be established cleanly. Check credentials or remote connectivity.</p>
|
||||
</div>"""
|
||||
return render_page(title="Traffic Control", body_html=body)
|
||||
|
||||
runnable_count = len(snapshot.runnable)
|
||||
leased_count = len(snapshot.leased)
|
||||
blocked_count = len(snapshot.blocked)
|
||||
controller_count = len(snapshot.needs_controller)
|
||||
terminal_count = len(snapshot.terminal_complete)
|
||||
|
||||
summary_bar = f"""<div class="health-card" style="display:flex; flex-wrap:wrap; gap:1rem; align-items:center;">
|
||||
<div><strong>Runnable:</strong> <span class="badge badge-health-ok">{runnable_count}</span></div>
|
||||
<div><strong>Leased:</strong> <span class="badge badge-claimed">{leased_count}</span></div>
|
||||
<div><strong>Blocked:</strong> <span class="badge badge-blocked">{blocked_count}</span></div>
|
||||
<div><strong>Needs Controller:</strong> <span class="badge badge-duplicate">{controller_count}</span></div>
|
||||
<div><strong>Terminal Complete:</strong> <span class="badge badge-stale">{terminal_count}</span></div>
|
||||
</div>"""
|
||||
|
||||
next_roles_html = _render_next_roles(snapshot.next_roles)
|
||||
|
||||
sections_html = f"""
|
||||
<div class="prompt-card">
|
||||
<h3>1. Runnable Lanes (Ready for Allocation)</h3>
|
||||
<p class="muted">Safe work items with no unmet dependencies or active leases. Safe for allocation.</p>
|
||||
{_render_traffic_table(snapshot.runnable, "No runnable items ready for allocation.")}
|
||||
</div>
|
||||
|
||||
<div class="prompt-card">
|
||||
<h3>2. In-Progress Work (Active Leases)</h3>
|
||||
<p class="muted">Work items currently leased and actively being worked by an assigned role session.</p>
|
||||
{_render_traffic_table(snapshot.leased, "No active leases in flight.")}
|
||||
</div>
|
||||
|
||||
<div class="prompt-card">
|
||||
<h3>3. Blocked Items (Dependencies / Locks)</h3>
|
||||
<p class="muted">Items blocked by unmet dependency issues, status:blocked, or active terminal review locks. Never presented as safe.</p>
|
||||
{_render_traffic_table(snapshot.blocked, "No blocked items.")}
|
||||
</div>
|
||||
|
||||
<div class="prompt-card">
|
||||
<h3>4. Needs Controller Intervention</h3>
|
||||
<p class="muted">Items requiring controller routing, diagnosis, or cross-role assignment.</p>
|
||||
{_render_traffic_table(snapshot.needs_controller, "No items requiring controller intervention.")}
|
||||
</div>
|
||||
|
||||
<div class="prompt-card">
|
||||
<h3>5. Terminal / Complete Candidates</h3>
|
||||
<p class="muted">Items ready for terminal reconciliation or post-merge worktree cleanup.</p>
|
||||
{_render_traffic_table(snapshot.terminal_complete, "No terminal complete candidates.")}
|
||||
</div>
|
||||
"""
|
||||
|
||||
body = f"""<h2>Workflow Traffic Control</h2>
|
||||
<p class="meta">Repository: <code>{escape(snapshot.repo_label)}</code></p>
|
||||
{summary_bar}
|
||||
{next_roles_html}
|
||||
{sections_html}"""
|
||||
|
||||
return render_page(title="Traffic Control", body_html=body)
|
||||
Reference in New Issue
Block a user