Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2068bae341 | ||
|
|
9517834913 | ||
|
|
824c42f7e3 | ||
|
|
578c44b685 | ||
|
|
9f686253eb | ||
|
|
1cbbde0089 |
@@ -1,53 +0,0 @@
|
|||||||
# MCP restart classes and blast-radius permissions (#663)
|
|
||||||
|
|
||||||
This is the machine-enforced class matrix used by
|
|
||||||
`restart_coordinator.RESTART_CLASS_POLICIES`. It implements the narrower-first
|
|
||||||
recovery ladder from #655 and the authorization policy from #656, using the
|
|
||||||
path inventory from #657 and the impact coordinator from #658. Product and
|
|
||||||
delivery lineage: vision #652 and roadmap #653.
|
|
||||||
|
|
||||||
Unknown class names are denied. The coordinator requires both the class
|
|
||||||
permission and an eligible request role. Approval gates are additional: a
|
|
||||||
caller cannot turn a request permission into execution authority.
|
|
||||||
|
|
||||||
| Restart class | Required permission | Expected blast radius | Drain requirement | Approval requirement | Audit requirement | Recovery behavior |
|
|
||||||
|---|---|---|---|---|---|---|
|
|
||||||
| `client_reconnect` | `mcp.reconnect.client` | none | none | self service | class, actor, client namespace, reason, outcome | Reconnect only the caller's client transport. No daemon or peer work changes. |
|
|
||||||
| `session_reconnect` | `mcp.reconnect.session` | low | requesting-session safe point | self service | class, actor, session, reason, outcome | Rebind identity, capability, and workspace state for one session. |
|
|
||||||
| `worker_restart` | `mcp.restart.worker.request` | low | target worker | controller approval + automated gates | class, actor, worker, approval, scoped drain, outcome | Restart one worker after its own leases and mutations drain. |
|
|
||||||
| `role_runtime_restart` | `mcp.restart.role_runtime.request` | medium | target role runtime | controller approval + automated gates | class, actor, role namespace, approval, scoped drain, outcome | Restart and re-probe one role runtime; unrelated roles remain available. |
|
|
||||||
| `connector_restart` | `mcp.restart.connector.request` | medium | target connector | controller approval + automated gates | class, actor, connector, approval, scoped drain, outcome | Restart one connector while unrelated runtimes remain available. |
|
|
||||||
| `configuration_reload` | `mcp.reload.configuration.request` | low | mutation quiesce | controller approval + automated gates | class, actor, configuration revision, approval, outcome | Gracefully reload configuration without replacing the daemon. |
|
|
||||||
| `rolling_mcp_restart` | `mcp.restart.rolling.request` | medium | one instance at a time | controller approval + automated gates | class, actor, instance order, approval, per-instance drains, outcome | Drain, restart, verify, and restore each instance before advancing. |
|
|
||||||
| `full_mcp_restart` | `mcp.restart.full.request` | high | all sessions and mutations | controller approval + automated gates | class, actor, full impact, approval, full drain proof, outcome | Replace the complete MCP runtime only after a verified full drain. |
|
|
||||||
| `host_restart` | `mcp.restart.host.request` | high | all host work | controller approval + infrastructure operator | class, actor, host/change or incident id, approval, full drain proof, outcome | Hand off to infrastructure ownership and reconcile every runtime afterward. |
|
|
||||||
|
|
||||||
## Drain boundary
|
|
||||||
|
|
||||||
Only `full_mcp_restart` and `host_restart` set `full_drain_required=true`.
|
|
||||||
Reconnects and configuration reloads do not disrupt peer sessions. Worker,
|
|
||||||
role-runtime, and connector restarts evaluate only their explicitly named
|
|
||||||
target. Rolling restart drains one instance at a time. Missing required target
|
|
||||||
scope denies the request rather than silently widening it to a full restart.
|
|
||||||
|
|
||||||
## Permission and approval boundary
|
|
||||||
|
|
||||||
Author, reviewer, merger, and reconciler roles may self-request reconnects and
|
|
||||||
request scoped worker/role/connector/reload recovery. They cannot request
|
|
||||||
rolling, full, or host restart classes. Controller/operator/admin roles may
|
|
||||||
request the broader classes, while execution remains operator/admin-owned.
|
|
||||||
Controller approval is independently required for every class above a session
|
|
||||||
reconnect. Host restart additionally requires infrastructure-operator proof.
|
|
||||||
|
|
||||||
The MCP request tool derives class permissions from its authenticated runtime
|
|
||||||
role. It does not accept caller-supplied permissions. Controller and operator
|
|
||||||
authorization are read from the already-running daemon environment, never
|
|
||||||
from a request argument.
|
|
||||||
|
|
||||||
## Audit and failure behavior
|
|
||||||
|
|
||||||
Every impact audit and every console restart/reload audit includes a
|
|
||||||
`restart_class` field. The impact audit also includes the exact
|
|
||||||
`required_permission`. Unknown classes, missing permissions, ineligible roles,
|
|
||||||
missing approval, missing scoped targets, and incomplete inventory all deny
|
|
||||||
fail closed. Manual process kills remain forbidden and contaminating (#630).
|
|
||||||
@@ -12,12 +12,6 @@ restart/reload/kill paths). The **mutative apply** path — actually performing
|
|||||||
restart — is a later child gated by a drain proof and is explicitly out of
|
restart — is a later child gated by a drain proof and is explicitly out of
|
||||||
scope here.
|
scope here.
|
||||||
|
|
||||||
The coordinator now routes every request through the restart-class policy
|
|
||||||
matrix defined for #663. See
|
|
||||||
[`mcp-restart-classes.md`](./mcp-restart-classes.md) for permissions, expected
|
|
||||||
blast radius, scoped drain and approval requirements, audit fields, and
|
|
||||||
recovery behavior for all nine classes.
|
|
||||||
|
|
||||||
## Components
|
## Components
|
||||||
|
|
||||||
| Piece | Where | Responsibility |
|
| Piece | Where | Responsibility |
|
||||||
@@ -83,10 +77,7 @@ authorization is present.
|
|||||||
```text
|
```text
|
||||||
gitea_request_mcp_restart(remote, host, org, repo,
|
gitea_request_mcp_restart(remote, host, org, repo,
|
||||||
dry_run=True, request_override=False,
|
dry_run=True, request_override=False,
|
||||||
session_id=None, limit=200,
|
session_id=None, limit=200)
|
||||||
restart_class="full_mcp_restart",
|
|
||||||
target_session_id=None, target_role=None,
|
|
||||||
target_connector=None)
|
|
||||||
```
|
```
|
||||||
|
|
||||||
Read-only, dry-run, and it **never restarts anything**. `apply_supported` is
|
Read-only, dry-run, and it **never restarts anything**. `apply_supported` is
|
||||||
@@ -96,8 +87,7 @@ apply is gated by a drain proof (a separate child).
|
|||||||
## Audit
|
## Audit
|
||||||
|
|
||||||
Every evaluation carries an `audit_record` (event, coordinator version, verdict,
|
Every evaluation carries an `audit_record` (event, coordinator version, verdict,
|
||||||
restart class, required permission, allow decision, blast radius, counts,
|
allow decision, blast radius, counts, timestamp) so restart decisions are
|
||||||
timestamp) so restart decisions are
|
|
||||||
auditable. No secrets flow through the coordinator — session ids, pids, and
|
auditable. No secrets flow through the coordinator — session ids, pids, and
|
||||||
profiles are operational metadata only.
|
profiles are operational metadata only.
|
||||||
|
|
||||||
|
|||||||
+831
@@ -0,0 +1,831 @@
|
|||||||
|
"""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"
|
||||||
|
|
||||||
|
# The only values accepted as an explicit acknowledgement from a live session.
|
||||||
|
# Anything else — including a missing entry, a null, or a "pending"/"stale"
|
||||||
|
# marker — leaves that session unacknowledged.
|
||||||
|
_ACK_TOKENS = frozenset({"ack", "acked", "acknowledged"})
|
||||||
|
|
||||||
|
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 _live_session_count(impact_report: Mapping[str, Any]) -> int | None:
|
||||||
|
"""Other-live-session count from the report, or ``None`` when unproven.
|
||||||
|
|
||||||
|
Only a real, non-negative integer counts. A missing ``counts`` block, a
|
||||||
|
malformed one, a non-integer, a bool (``True`` is an ``int`` in Python), or
|
||||||
|
a negative value all return ``None`` so the caller fails closed rather than
|
||||||
|
treating an unreadable report as "nobody was live".
|
||||||
|
"""
|
||||||
|
|
||||||
|
counts = impact_report.get("counts")
|
||||||
|
if not isinstance(counts, Mapping):
|
||||||
|
return None
|
||||||
|
value = counts.get("sessions_live_other")
|
||||||
|
if isinstance(value, bool) or not isinstance(value, int) or value < 0:
|
||||||
|
return None
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def _is_acknowledged(value: Any) -> bool:
|
||||||
|
"""True only for an explicit acknowledgement token.
|
||||||
|
|
||||||
|
Deliberately strict: the value must already be a string carrying one of the
|
||||||
|
recognised tokens. Non-strings are not coerced, so ``None``, timestamps,
|
||||||
|
objects, and states such as "pending" or "stale" are never read as an
|
||||||
|
acknowledgement.
|
||||||
|
"""
|
||||||
|
|
||||||
|
if not isinstance(value, str):
|
||||||
|
return False
|
||||||
|
return value.strip().lower() in _ACK_TOKENS
|
||||||
|
|
||||||
|
|
||||||
|
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.
|
||||||
|
#
|
||||||
|
# Whether acknowledgements are *required* is answered by the impact report,
|
||||||
|
# never inferred from the shape of the drain state. Previously an absent
|
||||||
|
# ``acks`` key collapsed to ``{}`` and was read as "no session needed to
|
||||||
|
# acknowledge", so a proof minted clean — and the restart gate allowed —
|
||||||
|
# while the report still showed other live sessions. Absence of evidence is
|
||||||
|
# not evidence of absence: missing, malformed, stale, or otherwise unproven
|
||||||
|
# acknowledgement data fails closed whenever live sessions require
|
||||||
|
# acknowledgement, and only explicitly verified acknowledgement evidence
|
||||||
|
# may permit the operation.
|
||||||
|
sessions_live_other = _live_session_count(impact_report)
|
||||||
|
live_count_known = sessions_live_other is not None
|
||||||
|
# An unknown or malformed count fails closed: it cannot prove nobody had to
|
||||||
|
# acknowledge, so acknowledgement stays required.
|
||||||
|
acks_required = (not live_count_known) or sessions_live_other > 0
|
||||||
|
|
||||||
|
raw_acks = drain_state.get("acks")
|
||||||
|
acks_present = isinstance(raw_acks, Mapping)
|
||||||
|
ack_values = list(raw_acks.values()) if acks_present else []
|
||||||
|
acked_count = sum(1 for value in ack_values if _is_acknowledged(value))
|
||||||
|
every_entry_acked = bool(ack_values) and acked_count == len(ack_values)
|
||||||
|
covers_live_sessions = live_count_known and acked_count >= sessions_live_other
|
||||||
|
# Strict by construction: only an explicit ``True`` counts, so an absent,
|
||||||
|
# null, or non-boolean ``ack_timeout_policy_applied`` can never open this
|
||||||
|
# gate on its own.
|
||||||
|
timeout_policy = _bool_input(drain_state.get("ack_timeout_policy_applied"))
|
||||||
|
|
||||||
|
# Present-but-unacknowledged entries always fail closed, even when the
|
||||||
|
# report claims nobody was live: the drain state naming an outstanding
|
||||||
|
# session contradicts that claim, and the safe reading of a contradiction
|
||||||
|
# is that an acknowledgement is still owed.
|
||||||
|
outstanding_entries = acks_present and bool(ack_values) and not every_entry_acked
|
||||||
|
|
||||||
|
if timeout_policy and (acks_required or outstanding_entries):
|
||||||
|
acks_ok = True
|
||||||
|
detail = "explicit ack timeout policy applied and recorded"
|
||||||
|
elif outstanding_entries:
|
||||||
|
acks_ok = False
|
||||||
|
detail = (
|
||||||
|
f"{len(ack_values) - acked_count} of {len(ack_values)} "
|
||||||
|
"acknowledgement entries are unparseable, stale, or not "
|
||||||
|
"acknowledged (fail closed)"
|
||||||
|
)
|
||||||
|
elif not acks_required:
|
||||||
|
acks_ok = True
|
||||||
|
detail = (
|
||||||
|
"impact report proves no other live sessions required to "
|
||||||
|
"acknowledge (sessions_live_other=0)"
|
||||||
|
)
|
||||||
|
elif every_entry_acked and covers_live_sessions:
|
||||||
|
acks_ok = True
|
||||||
|
detail = (
|
||||||
|
f"{acked_count} acknowledgement entr"
|
||||||
|
f"{'y' if acked_count == 1 else 'ies'} verified; covers "
|
||||||
|
f"sessions_live_other={sessions_live_other}"
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
acks_ok = False
|
||||||
|
if not live_count_known:
|
||||||
|
detail = (
|
||||||
|
"impact report does not prove the live-session count; "
|
||||||
|
"acknowledgement required and unproven (fail closed)"
|
||||||
|
)
|
||||||
|
elif raw_acks is None:
|
||||||
|
detail = (
|
||||||
|
"acknowledgement evidence absent while "
|
||||||
|
f"sessions_live_other={sessions_live_other} require "
|
||||||
|
"acknowledgement (fail closed)"
|
||||||
|
)
|
||||||
|
elif not acks_present:
|
||||||
|
detail = (
|
||||||
|
"acknowledgement evidence malformed: expected a mapping, got "
|
||||||
|
f"{type(raw_acks).__name__} (fail closed)"
|
||||||
|
)
|
||||||
|
elif not ack_values:
|
||||||
|
detail = (
|
||||||
|
"acknowledgement mapping empty while "
|
||||||
|
f"sessions_live_other={sessions_live_other} require "
|
||||||
|
"acknowledgement (fail closed)"
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
detail = (
|
||||||
|
f"acknowledgements cover only {acked_count} session(s) but the "
|
||||||
|
f"report shows sessions_live_other={sessions_live_other} "
|
||||||
|
"(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,
|
||||||
|
)
|
||||||
+59
-37
@@ -2068,6 +2068,7 @@ import lease_lifecycle # noqa: E402
|
|||||||
import lease_policy # noqa: E402
|
import lease_policy # noqa: E402
|
||||||
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
|
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
|
||||||
import restart_coordinator # noqa: E402 # #658 MCP restart coordinator/impact
|
import restart_coordinator # noqa: E402 # #658 MCP restart coordinator/impact
|
||||||
|
import drain_proof # noqa: E402 # #661 pre-restart drain proof and hard gate
|
||||||
import incident_bridge # noqa: E402
|
import incident_bridge # noqa: E402
|
||||||
import sentry_observability # noqa: E402 (#606 optional Sentry observability)
|
import sentry_observability # noqa: E402 (#606 optional Sentry observability)
|
||||||
import sentry_incident_bridge # noqa: E402 (#607 Sentry→Gitea incident bridge)
|
import sentry_incident_bridge # noqa: E402 (#607 Sentry→Gitea incident bridge)
|
||||||
@@ -22342,24 +22343,26 @@ def gitea_request_mcp_restart(
|
|||||||
request_override: bool = False,
|
request_override: bool = False,
|
||||||
session_id: str | None = None,
|
session_id: str | None = None,
|
||||||
limit: int = 200,
|
limit: int = 200,
|
||||||
restart_class: str = "full_mcp_restart",
|
drain_proof_json: str | None = None,
|
||||||
target_session_id: str | None = None,
|
request_break_glass: bool = False,
|
||||||
target_role: str | None = None,
|
|
||||||
target_connector: str | None = None,
|
|
||||||
) -> dict:
|
) -> dict:
|
||||||
"""Evaluate a proposed MCP restart and return an impact preview (#658).
|
"""Evaluate a proposed MCP restart and return an impact preview (#658).
|
||||||
|
|
||||||
Central restart coordinator: resolves the requested restart class, gathers
|
Central restart coordinator: gathers live control-plane state (sessions,
|
||||||
live control-plane state (sessions,
|
|
||||||
leases/locks, in-flight issue/PR work, mutations, worktrees) and returns a
|
leases/locks, in-flight issue/PR work, mutations, worktrees) and returns a
|
||||||
blast-radius impact report with a ``safe`` / ``unsafe`` / ``override``
|
blast-radius impact report with a ``safe`` / ``unsafe`` / ``override``
|
||||||
verdict, so the console (#642/#652) and operators can see what a restart
|
verdict, so the console (#642/#652) and operators can see what a restart
|
||||||
would disrupt *before* any concurrent LLM work is destroyed.
|
would disrupt *before* any concurrent LLM work is destroyed.
|
||||||
|
|
||||||
This tool is **dry-run and never restarts anything.** The mutative apply
|
This tool **never restarts a process.** In dry-run (the default) it returns
|
||||||
path is a separate child gated by a drain proof (non-goal here); calling
|
only the impact preview. With ``dry_run=False`` it enforces the #661 hard
|
||||||
with ``dry_run=False`` still performs no restart and reports that apply is
|
gate: the apply request must present a valid, unexpired, clean drain proof
|
||||||
not yet available.
|
(``drain_proof_json``) or it is denied and a durable incident descriptor is
|
||||||
|
returned under ``incident``. Break-glass is the only bypass and is honoured
|
||||||
|
only when ``request_break_glass`` is set *and* the environment carries
|
||||||
|
``GITEA_BREAKGLASS_RESTART_AUTHORIZATION``. Even an authorized gate performs
|
||||||
|
no restart here; actual execution is a further child. The gate outcome is
|
||||||
|
reported under ``apply_gate`` / ``apply_authorized``.
|
||||||
|
|
||||||
Operator override authority is read from the process environment
|
Operator override authority is read from the process environment
|
||||||
(``GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION``), never self-asserted by
|
(``GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION``), never self-asserted by
|
||||||
@@ -22444,9 +22447,6 @@ def gitea_request_mcp_restart(
|
|||||||
|
|
||||||
profile = get_profile()
|
profile = get_profile()
|
||||||
profile_name = (profile.get("profile_name") or "").strip() or "session"
|
profile_name = (profile.get("profile_name") or "").strip() or "session"
|
||||||
requester_role = (
|
|
||||||
profile.get("role_kind") or profile.get("role") or ""
|
|
||||||
).strip().lower()
|
|
||||||
sid = (session_id or "").strip() or f"{profile_name}-{os.getpid()}"
|
sid = (session_id or "").strip() or f"{profile_name}-{os.getpid()}"
|
||||||
|
|
||||||
# Override authority is read from the environment only — a worker session
|
# Override authority is read from the environment only — a worker session
|
||||||
@@ -22456,15 +22456,6 @@ def gitea_request_mcp_restart(
|
|||||||
(os.environ.get("GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION") or "").strip()
|
(os.environ.get("GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION") or "").strip()
|
||||||
)
|
)
|
||||||
operator_override = bool(request_override and operator_authorized)
|
operator_override = bool(request_override and operator_authorized)
|
||||||
controller_approved = bool(
|
|
||||||
(
|
|
||||||
os.environ.get("GITEA_CONTROLLER_RESTART_APPROVAL_AUTHORIZATION")
|
|
||||||
or ""
|
|
||||||
).strip()
|
|
||||||
)
|
|
||||||
requester_permissions = restart_coordinator.permissions_for_role(
|
|
||||||
requester_role
|
|
||||||
)
|
|
||||||
|
|
||||||
inventory = {
|
inventory = {
|
||||||
"sessions": sessions,
|
"sessions": sessions,
|
||||||
@@ -22479,14 +22470,6 @@ def gitea_request_mcp_restart(
|
|||||||
operator_override=operator_override,
|
operator_override=operator_override,
|
||||||
requesting_session_id=sid,
|
requesting_session_id=sid,
|
||||||
dry_run=True, # coordinator is always analysis-only (#658)
|
dry_run=True, # coordinator is always analysis-only (#658)
|
||||||
restart_class=restart_class,
|
|
||||||
requester_role=requester_role,
|
|
||||||
requester_permissions=requester_permissions,
|
|
||||||
controller_approved=controller_approved,
|
|
||||||
operator_authorized=operator_authorized,
|
|
||||||
target_session_id=target_session_id,
|
|
||||||
target_role=target_role,
|
|
||||||
target_connector=target_connector,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
payload = report.as_dict()
|
payload = report.as_dict()
|
||||||
@@ -22498,15 +22481,54 @@ def gitea_request_mcp_restart(
|
|||||||
payload["requesting_session_id"] = sid
|
payload["requesting_session_id"] = sid
|
||||||
payload["operator_override_requested"] = bool(request_override)
|
payload["operator_override_requested"] = bool(request_override)
|
||||||
payload["operator_override_authorized"] = operator_authorized
|
payload["operator_override_authorized"] = operator_authorized
|
||||||
payload["controller_approval_authorized"] = controller_approved
|
# Actual restart execution remains a further child; this tool never restarts
|
||||||
payload["requester_role"] = requester_role
|
# a process. What #661 adds is the *hard gate*: an apply request (dry_run
|
||||||
payload["requester_permissions"] = list(requester_permissions)
|
# False) must present a valid, unexpired, clean drain proof, or it is denied
|
||||||
|
# and a durable incident is raised. Break-glass is the only bypass and its
|
||||||
|
# authorization is read from the environment, never self-asserted.
|
||||||
payload["apply_supported"] = False
|
payload["apply_supported"] = False
|
||||||
if not dry_run:
|
if not dry_run:
|
||||||
payload["reasons"] = list(payload.get("reasons") or []) + [
|
proof_obj: dict | None = None
|
||||||
"apply requested but not supported: sanctioned restart apply is "
|
proof_parse_error: str | None = None
|
||||||
"gated by a drain proof (separate child); no restart performed (#658)"
|
if drain_proof_json:
|
||||||
]
|
try:
|
||||||
|
parsed = json.loads(drain_proof_json)
|
||||||
|
proof_obj = parsed if isinstance(parsed, dict) else None
|
||||||
|
if proof_obj is None:
|
||||||
|
proof_parse_error = "drain_proof_json is not a JSON object"
|
||||||
|
except (ValueError, TypeError) as exc:
|
||||||
|
proof_parse_error = f"invalid drain_proof_json: {_redact(str(exc))}"
|
||||||
|
|
||||||
|
break_glass_authorized = bool(
|
||||||
|
(
|
||||||
|
os.environ.get("GITEA_BREAKGLASS_RESTART_AUTHORIZATION") or ""
|
||||||
|
).strip()
|
||||||
|
)
|
||||||
|
break_glass = bool(request_break_glass and break_glass_authorized)
|
||||||
|
|
||||||
|
expected_fp = drain_proof.impact_fingerprint(report.as_dict())
|
||||||
|
gate = drain_proof.gate_apply_restart(
|
||||||
|
proof=proof_obj,
|
||||||
|
break_glass=break_glass,
|
||||||
|
expected_impact_fingerprint=expected_fp,
|
||||||
|
requesting_session_id=sid,
|
||||||
|
)
|
||||||
|
gate_payload = gate.as_dict()
|
||||||
|
if proof_parse_error and not break_glass:
|
||||||
|
gate_payload["reasons"] = [proof_parse_error] + list(
|
||||||
|
gate_payload.get("reasons") or []
|
||||||
|
)
|
||||||
|
payload["apply_gate"] = gate_payload
|
||||||
|
payload["apply_authorized"] = gate.allow
|
||||||
|
payload["break_glass_requested"] = bool(request_break_glass)
|
||||||
|
payload["break_glass_authorized"] = break_glass_authorized
|
||||||
|
# Even an authorized gate performs no restart here: execution is a later
|
||||||
|
# child. The gate proves the apply path *would* be permitted.
|
||||||
|
payload["reasons"] = list(payload.get("reasons") or []) + list(
|
||||||
|
gate_payload.get("reasons") or []
|
||||||
|
)
|
||||||
|
if not gate.allow and gate.incident is not None:
|
||||||
|
payload["incident"] = gate.incident
|
||||||
return payload
|
return payload
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
+14
-370
@@ -28,12 +28,11 @@ from __future__ import annotations
|
|||||||
|
|
||||||
from dataclasses import dataclass, field
|
from dataclasses import dataclass, field
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from enum import Enum
|
|
||||||
from typing import Any, Mapping, Sequence
|
from typing import Any, Mapping, Sequence
|
||||||
|
|
||||||
import lease_lifecycle
|
import lease_lifecycle
|
||||||
|
|
||||||
COORDINATOR_VERSION = "1.1.0-issue-663"
|
COORDINATOR_VERSION = "1.0.0-issue-658"
|
||||||
|
|
||||||
# Restart verdicts. Exactly the three the acceptance criteria name.
|
# Restart verdicts. Exactly the three the acceptance criteria name.
|
||||||
VERDICT_SAFE = "safe"
|
VERDICT_SAFE = "safe"
|
||||||
@@ -55,194 +54,6 @@ LEASE_FRESHNESS_LIVE = "active"
|
|||||||
DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS = 900
|
DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS = 900
|
||||||
|
|
||||||
|
|
||||||
class RestartClass(str, Enum):
|
|
||||||
"""The only restart/recovery classes accepted by the coordinator."""
|
|
||||||
|
|
||||||
CLIENT_RECONNECT = "client_reconnect"
|
|
||||||
SESSION_RECONNECT = "session_reconnect"
|
|
||||||
WORKER_RESTART = "worker_restart"
|
|
||||||
ROLE_RUNTIME_RESTART = "role_runtime_restart"
|
|
||||||
CONNECTOR_RESTART = "connector_restart"
|
|
||||||
CONFIGURATION_RELOAD = "configuration_reload"
|
|
||||||
ROLLING_MCP_RESTART = "rolling_mcp_restart"
|
|
||||||
FULL_MCP_RESTART = "full_mcp_restart"
|
|
||||||
HOST_RESTART = "host_restart"
|
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True)
|
|
||||||
class RestartClassPolicy:
|
|
||||||
"""Least-privilege policy for one :class:`RestartClass`."""
|
|
||||||
|
|
||||||
restart_class: RestartClass
|
|
||||||
required_permission: str
|
|
||||||
expected_blast_radius: str
|
|
||||||
drain_requirement: str
|
|
||||||
full_drain_required: bool
|
|
||||||
approval_requirement: str
|
|
||||||
audit_requirement: str
|
|
||||||
recovery_behavior: str
|
|
||||||
request_roles: tuple[str, ...]
|
|
||||||
execution_roles: tuple[str, ...]
|
|
||||||
|
|
||||||
def as_dict(self) -> dict[str, Any]:
|
|
||||||
return {
|
|
||||||
"restart_class": self.restart_class.value,
|
|
||||||
"required_permission": self.required_permission,
|
|
||||||
"expected_blast_radius": self.expected_blast_radius,
|
|
||||||
"drain_requirement": self.drain_requirement,
|
|
||||||
"full_drain_required": self.full_drain_required,
|
|
||||||
"approval_requirement": self.approval_requirement,
|
|
||||||
"audit_requirement": self.audit_requirement,
|
|
||||||
"recovery_behavior": self.recovery_behavior,
|
|
||||||
"request_roles": list(self.request_roles),
|
|
||||||
"execution_roles": list(self.execution_roles),
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
WORKER_ROLES = ("author", "reviewer", "merger", "reconciler")
|
|
||||||
CONTROL_ROLES = ("controller", "operator", "admin")
|
|
||||||
ALL_REQUEST_ROLES = WORKER_ROLES + CONTROL_ROLES
|
|
||||||
|
|
||||||
RESTART_CLASS_POLICIES: dict[RestartClass, RestartClassPolicy] = {
|
|
||||||
RestartClass.CLIENT_RECONNECT: RestartClassPolicy(
|
|
||||||
RestartClass.CLIENT_RECONNECT,
|
|
||||||
"mcp.reconnect.client",
|
|
||||||
BLAST_NONE,
|
|
||||||
"none",
|
|
||||||
False,
|
|
||||||
"self_service",
|
|
||||||
"record class, actor, client namespace, reason, and outcome",
|
|
||||||
"Reconnect only the caller's client transport; no daemon or peer session changes.",
|
|
||||||
ALL_REQUEST_ROLES,
|
|
||||||
ALL_REQUEST_ROLES,
|
|
||||||
),
|
|
||||||
RestartClass.SESSION_RECONNECT: RestartClassPolicy(
|
|
||||||
RestartClass.SESSION_RECONNECT,
|
|
||||||
"mcp.reconnect.session",
|
|
||||||
BLAST_LOW,
|
|
||||||
"requesting_session_safe_point",
|
|
||||||
False,
|
|
||||||
"self_service",
|
|
||||||
"record class, actor, session id, reason, and outcome",
|
|
||||||
"Rebind identity, capability, and workspace state for one session.",
|
|
||||||
ALL_REQUEST_ROLES,
|
|
||||||
ALL_REQUEST_ROLES,
|
|
||||||
),
|
|
||||||
RestartClass.WORKER_RESTART: RestartClassPolicy(
|
|
||||||
RestartClass.WORKER_RESTART,
|
|
||||||
"mcp.restart.worker.request",
|
|
||||||
BLAST_LOW,
|
|
||||||
"target_worker",
|
|
||||||
False,
|
|
||||||
"controller_approval_and_automated_gates",
|
|
||||||
"record class, actor, target worker, approval, drain proof, and outcome",
|
|
||||||
"Restart one worker after its own lease and mutation scope is drained.",
|
|
||||||
ALL_REQUEST_ROLES,
|
|
||||||
("operator", "admin"),
|
|
||||||
),
|
|
||||||
RestartClass.ROLE_RUNTIME_RESTART: RestartClassPolicy(
|
|
||||||
RestartClass.ROLE_RUNTIME_RESTART,
|
|
||||||
"mcp.restart.role_runtime.request",
|
|
||||||
BLAST_MEDIUM,
|
|
||||||
"target_role_runtime",
|
|
||||||
False,
|
|
||||||
"controller_approval_and_automated_gates",
|
|
||||||
"record class, actor, role namespace, approval, drain proof, and outcome",
|
|
||||||
"Restart only the selected role runtime and then re-probe that namespace.",
|
|
||||||
ALL_REQUEST_ROLES,
|
|
||||||
("operator", "admin"),
|
|
||||||
),
|
|
||||||
RestartClass.CONNECTOR_RESTART: RestartClassPolicy(
|
|
||||||
RestartClass.CONNECTOR_RESTART,
|
|
||||||
"mcp.restart.connector.request",
|
|
||||||
BLAST_MEDIUM,
|
|
||||||
"target_connector",
|
|
||||||
False,
|
|
||||||
"controller_approval_and_automated_gates",
|
|
||||||
"record class, actor, connector id, approval, drain proof, and outcome",
|
|
||||||
"Restart one connector while unrelated role runtimes remain available.",
|
|
||||||
ALL_REQUEST_ROLES,
|
|
||||||
("operator", "admin"),
|
|
||||||
),
|
|
||||||
RestartClass.CONFIGURATION_RELOAD: RestartClassPolicy(
|
|
||||||
RestartClass.CONFIGURATION_RELOAD,
|
|
||||||
"mcp.reload.configuration.request",
|
|
||||||
BLAST_LOW,
|
|
||||||
"mutation_quiesce",
|
|
||||||
False,
|
|
||||||
"controller_approval_and_automated_gates",
|
|
||||||
"record class, actor, configuration revision, approval, and outcome",
|
|
||||||
"Gracefully reload configuration without replacing the daemon process.",
|
|
||||||
ALL_REQUEST_ROLES,
|
|
||||||
("operator", "admin"),
|
|
||||||
),
|
|
||||||
RestartClass.ROLLING_MCP_RESTART: RestartClassPolicy(
|
|
||||||
RestartClass.ROLLING_MCP_RESTART,
|
|
||||||
"mcp.restart.rolling.request",
|
|
||||||
BLAST_MEDIUM,
|
|
||||||
"one_instance_at_a_time",
|
|
||||||
False,
|
|
||||||
"controller_approval_and_automated_gates",
|
|
||||||
"record class, actor, instance order, approval, per-instance drains, and outcome",
|
|
||||||
"Drain, restart, verify, and restore one instance before advancing to the next.",
|
|
||||||
CONTROL_ROLES,
|
|
||||||
("operator", "admin"),
|
|
||||||
),
|
|
||||||
RestartClass.FULL_MCP_RESTART: RestartClassPolicy(
|
|
||||||
RestartClass.FULL_MCP_RESTART,
|
|
||||||
"mcp.restart.full.request",
|
|
||||||
BLAST_HIGH,
|
|
||||||
"all_sessions_and_mutations",
|
|
||||||
True,
|
|
||||||
"controller_approval_and_automated_gates",
|
|
||||||
"record class, actor, full impact report, approval, drain proof, and outcome",
|
|
||||||
"Stop and restore the complete MCP runtime only after a verified full drain.",
|
|
||||||
CONTROL_ROLES,
|
|
||||||
("operator", "admin"),
|
|
||||||
),
|
|
||||||
RestartClass.HOST_RESTART: RestartClassPolicy(
|
|
||||||
RestartClass.HOST_RESTART,
|
|
||||||
"mcp.restart.host.request",
|
|
||||||
BLAST_HIGH,
|
|
||||||
"all_host_work",
|
|
||||||
True,
|
|
||||||
"controller_approval_plus_infrastructure_operator",
|
|
||||||
"record class, actor, host, incident or change id, approval, drain proof, and outcome",
|
|
||||||
"Hand off to infrastructure ownership; reconcile every runtime after the host returns.",
|
|
||||||
("controller", "operator", "admin"),
|
|
||||||
("operator", "admin"),
|
|
||||||
),
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
def resolve_restart_class(value: RestartClass | str) -> RestartClass:
|
|
||||||
"""Resolve a restart class or fail closed for an unknown value."""
|
|
||||||
|
|
||||||
if isinstance(value, RestartClass):
|
|
||||||
return value
|
|
||||||
try:
|
|
||||||
return RestartClass(str(value).strip())
|
|
||||||
except ValueError as exc:
|
|
||||||
raise ValueError(f"unknown restart class {value!r}; deny (fail closed)") from exc
|
|
||||||
|
|
||||||
|
|
||||||
def restart_class_policy(value: RestartClass | str) -> RestartClassPolicy:
|
|
||||||
"""Return the canonical policy for *value*."""
|
|
||||||
|
|
||||||
return RESTART_CLASS_POLICIES[resolve_restart_class(value)]
|
|
||||||
|
|
||||||
|
|
||||||
def permissions_for_role(role: str | None) -> tuple[str, ...]:
|
|
||||||
"""Return request permissions granted to a workflow role by this policy."""
|
|
||||||
|
|
||||||
normalized = str(role or "").strip().lower()
|
|
||||||
return tuple(
|
|
||||||
policy.required_permission
|
|
||||||
for policy in RESTART_CLASS_POLICIES.values()
|
|
||||||
if normalized in policy.request_roles
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def _utc_now() -> datetime:
|
def _utc_now() -> datetime:
|
||||||
return datetime.now(timezone.utc)
|
return datetime.now(timezone.utc)
|
||||||
|
|
||||||
@@ -264,7 +75,6 @@ class SessionImpact:
|
|||||||
heartbeat_stale: bool
|
heartbeat_stale: bool
|
||||||
is_requester: bool
|
is_requester: bool
|
||||||
live: bool
|
live: bool
|
||||||
connector: str | None = None
|
|
||||||
|
|
||||||
def as_dict(self) -> dict[str, Any]:
|
def as_dict(self) -> dict[str, Any]:
|
||||||
return {
|
return {
|
||||||
@@ -277,7 +87,6 @@ class SessionImpact:
|
|||||||
"heartbeat_stale": self.heartbeat_stale,
|
"heartbeat_stale": self.heartbeat_stale,
|
||||||
"is_requester": self.is_requester,
|
"is_requester": self.is_requester,
|
||||||
"live": self.live,
|
"live": self.live,
|
||||||
"connector": self.connector,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -296,7 +105,6 @@ class LeaseImpact:
|
|||||||
disruptive: bool
|
disruptive: bool
|
||||||
is_mutation: bool
|
is_mutation: bool
|
||||||
is_critical_section: bool
|
is_critical_section: bool
|
||||||
connector: str | None = None
|
|
||||||
|
|
||||||
def as_dict(self) -> dict[str, Any]:
|
def as_dict(self) -> dict[str, Any]:
|
||||||
return {
|
return {
|
||||||
@@ -311,7 +119,6 @@ class LeaseImpact:
|
|||||||
"disruptive": self.disruptive,
|
"disruptive": self.disruptive,
|
||||||
"is_mutation": self.is_mutation,
|
"is_mutation": self.is_mutation,
|
||||||
"is_critical_section": self.is_critical_section,
|
"is_critical_section": self.is_critical_section,
|
||||||
"connector": self.connector,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -320,13 +127,6 @@ class RestartImpactReport:
|
|||||||
"""Impact preview DTO returned to the console / operator (#642/#652)."""
|
"""Impact preview DTO returned to the console / operator (#642/#652)."""
|
||||||
|
|
||||||
coordinator_version: str
|
coordinator_version: str
|
||||||
restart_class: str
|
|
||||||
restart_policy: dict[str, Any]
|
|
||||||
policy_enforced: bool
|
|
||||||
permission_authorized: bool
|
|
||||||
role_authorized: bool
|
|
||||||
approval_satisfied: bool
|
|
||||||
authorization_reasons: list[str]
|
|
||||||
evaluated_at: str
|
evaluated_at: str
|
||||||
dry_run: bool
|
dry_run: bool
|
||||||
restart_performed: bool
|
restart_performed: bool
|
||||||
@@ -353,13 +153,6 @@ class RestartImpactReport:
|
|||||||
def as_dict(self) -> dict[str, Any]:
|
def as_dict(self) -> dict[str, Any]:
|
||||||
return {
|
return {
|
||||||
"coordinator_version": self.coordinator_version,
|
"coordinator_version": self.coordinator_version,
|
||||||
"restart_class": self.restart_class,
|
|
||||||
"restart_policy": dict(self.restart_policy),
|
|
||||||
"policy_enforced": self.policy_enforced,
|
|
||||||
"permission_authorized": self.permission_authorized,
|
|
||||||
"role_authorized": self.role_authorized,
|
|
||||||
"approval_satisfied": self.approval_satisfied,
|
|
||||||
"authorization_reasons": list(self.authorization_reasons),
|
|
||||||
"evaluated_at": self.evaluated_at,
|
"evaluated_at": self.evaluated_at,
|
||||||
"dry_run": self.dry_run,
|
"dry_run": self.dry_run,
|
||||||
"restart_performed": self.restart_performed,
|
"restart_performed": self.restart_performed,
|
||||||
@@ -413,7 +206,6 @@ def _classify_session(
|
|||||||
requesting_session_id and session_id == requesting_session_id
|
requesting_session_id and session_id == requesting_session_id
|
||||||
),
|
),
|
||||||
live=live,
|
live=live,
|
||||||
connector=(str(row.get("connector") or "").strip() or None),
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@@ -466,7 +258,6 @@ def _classify_lease(row: Mapping[str, Any]) -> LeaseImpact:
|
|||||||
disruptive=disruptive,
|
disruptive=disruptive,
|
||||||
is_mutation=is_mutation,
|
is_mutation=is_mutation,
|
||||||
is_critical_section=disruptive,
|
is_critical_section=disruptive,
|
||||||
connector=(str(row.get("connector") or "").strip() or None),
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@@ -488,14 +279,6 @@ def evaluate_restart_impact(
|
|||||||
requesting_session_id: str | None = None,
|
requesting_session_id: str | None = None,
|
||||||
dry_run: bool = True,
|
dry_run: bool = True,
|
||||||
session_heartbeat_stale_seconds: int = DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS,
|
session_heartbeat_stale_seconds: int = DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS,
|
||||||
restart_class: RestartClass | str | None = None,
|
|
||||||
requester_role: str | None = None,
|
|
||||||
requester_permissions: Sequence[str] | None = None,
|
|
||||||
controller_approved: bool = False,
|
|
||||||
operator_authorized: bool = False,
|
|
||||||
target_session_id: str | None = None,
|
|
||||||
target_role: str | None = None,
|
|
||||||
target_connector: str | None = None,
|
|
||||||
) -> RestartImpactReport:
|
) -> RestartImpactReport:
|
||||||
"""Evaluate a proposed MCP restart and return an impact preview.
|
"""Evaluate a proposed MCP restart and return an impact preview.
|
||||||
|
|
||||||
@@ -518,61 +301,6 @@ def evaluate_restart_impact(
|
|||||||
"""
|
"""
|
||||||
moment = now or _utc_now()
|
moment = now or _utc_now()
|
||||||
reasons: list[str] = []
|
reasons: list[str] = []
|
||||||
authorization_reasons: list[str] = []
|
|
||||||
|
|
||||||
# ``None`` preserves the pre-#663 impact-only API for callers that have not
|
|
||||||
# yet been migrated. All MCP requests pass an explicit class and therefore
|
|
||||||
# take the fail-closed policy path.
|
|
||||||
policy_enforced = restart_class is not None
|
|
||||||
try:
|
|
||||||
resolved_class = resolve_restart_class(
|
|
||||||
restart_class or RestartClass.FULL_MCP_RESTART
|
|
||||||
)
|
|
||||||
policy = RESTART_CLASS_POLICIES[resolved_class]
|
|
||||||
unknown_class = False
|
|
||||||
except ValueError as exc:
|
|
||||||
resolved_class = None
|
|
||||||
policy = None
|
|
||||||
unknown_class = True
|
|
||||||
authorization_reasons.append(str(exc))
|
|
||||||
|
|
||||||
normalized_role = str(requester_role or "").strip().lower()
|
|
||||||
granted = {str(p).strip() for p in (requester_permissions or ())}
|
|
||||||
if policy_enforced and policy is not None:
|
|
||||||
permission_authorized = policy.required_permission in granted
|
|
||||||
role_authorized = normalized_role in policy.request_roles
|
|
||||||
if not permission_authorized:
|
|
||||||
authorization_reasons.append(
|
|
||||||
f"missing required permission {policy.required_permission!r}"
|
|
||||||
)
|
|
||||||
if not role_authorized:
|
|
||||||
authorization_reasons.append(
|
|
||||||
f"role {normalized_role or 'unknown'!r} may not request "
|
|
||||||
f"{policy.restart_class.value}"
|
|
||||||
)
|
|
||||||
elif unknown_class:
|
|
||||||
permission_authorized = False
|
|
||||||
role_authorized = False
|
|
||||||
else:
|
|
||||||
permission_authorized = True
|
|
||||||
role_authorized = True
|
|
||||||
|
|
||||||
if policy_enforced and policy is not None:
|
|
||||||
approval = policy.approval_requirement
|
|
||||||
if approval == "self_service":
|
|
||||||
approval_satisfied = True
|
|
||||||
elif approval == "controller_approval_plus_infrastructure_operator":
|
|
||||||
approval_satisfied = bool(controller_approved and operator_authorized)
|
|
||||||
else:
|
|
||||||
approval_satisfied = bool(controller_approved)
|
|
||||||
if not approval_satisfied:
|
|
||||||
authorization_reasons.append(
|
|
||||||
f"approval requirement not satisfied: {approval}"
|
|
||||||
)
|
|
||||||
elif unknown_class:
|
|
||||||
approval_satisfied = False
|
|
||||||
else:
|
|
||||||
approval_satisfied = True
|
|
||||||
|
|
||||||
inventory_complete = bool(inventory.get("inventory_complete", False))
|
inventory_complete = bool(inventory.get("inventory_complete", False))
|
||||||
incomplete_reasons = [str(r) for r in (inventory.get("incomplete_reasons") or [])]
|
incomplete_reasons = [str(r) for r in (inventory.get("incomplete_reasons") or [])]
|
||||||
@@ -595,67 +323,15 @@ def evaluate_restart_impact(
|
|||||||
]
|
]
|
||||||
lease_impacts = [_classify_lease(l) for l in leases_raw]
|
lease_impacts = [_classify_lease(l) for l in leases_raw]
|
||||||
|
|
||||||
# Route impact through the selected class. Narrow classes never inherit a
|
# Only *other* live sessions and live leases constitute blast radius: a
|
||||||
# full-runtime drain merely because unrelated work exists.
|
# restart that would kill only the requesting session with no other work in
|
||||||
target_complete = True
|
# flight is safe.
|
||||||
if resolved_class in {
|
|
||||||
RestartClass.CLIENT_RECONNECT,
|
|
||||||
RestartClass.SESSION_RECONNECT,
|
|
||||||
RestartClass.CONFIGURATION_RELOAD,
|
|
||||||
}:
|
|
||||||
scoped_sessions: list[SessionImpact] = []
|
|
||||||
scoped_leases: list[LeaseImpact] = []
|
|
||||||
elif resolved_class == RestartClass.WORKER_RESTART:
|
|
||||||
selected_session = (target_session_id or "").strip()
|
|
||||||
target_complete = bool(selected_session)
|
|
||||||
scoped_sessions = [
|
|
||||||
s for s in session_impacts if s.session_id == selected_session
|
|
||||||
]
|
|
||||||
scoped_leases = [
|
|
||||||
l for l in lease_impacts if l.session_id == selected_session
|
|
||||||
]
|
|
||||||
elif resolved_class == RestartClass.ROLE_RUNTIME_RESTART:
|
|
||||||
selected_role = (target_role or "").strip().lower()
|
|
||||||
target_complete = bool(selected_role)
|
|
||||||
scoped_sessions = [
|
|
||||||
s for s in session_impacts if str(s.role or "").lower() == selected_role
|
|
||||||
]
|
|
||||||
scoped_leases = [
|
|
||||||
l for l in lease_impacts if str(l.role or "").lower() == selected_role
|
|
||||||
]
|
|
||||||
elif resolved_class == RestartClass.CONNECTOR_RESTART:
|
|
||||||
selected_connector = (target_connector or "").strip()
|
|
||||||
target_complete = bool(selected_connector)
|
|
||||||
scoped_sessions = [
|
|
||||||
s for s in session_impacts if s.connector == selected_connector
|
|
||||||
]
|
|
||||||
scoped_leases = [
|
|
||||||
l for l in lease_impacts if l.connector == selected_connector
|
|
||||||
]
|
|
||||||
else:
|
|
||||||
scoped_sessions = list(session_impacts)
|
|
||||||
scoped_leases = list(lease_impacts)
|
|
||||||
|
|
||||||
if policy_enforced and not target_complete:
|
|
||||||
authorization_reasons.append(
|
|
||||||
f"target required for {resolved_class.value if resolved_class else 'unknown class'}"
|
|
||||||
)
|
|
||||||
|
|
||||||
other_live_sessions = [
|
other_live_sessions = [
|
||||||
s for s in scoped_sessions if s.live and not s.is_requester
|
s for s in session_impacts if s.live and not s.is_requester
|
||||||
]
|
]
|
||||||
disruptive_leases = [l for l in scoped_leases if l.disruptive]
|
disruptive_leases = [l for l in lease_impacts if l.disruptive]
|
||||||
critical_sections = [l for l in scoped_leases if l.is_critical_section]
|
critical_sections = [l for l in lease_impacts if l.is_critical_section]
|
||||||
mutations = [l for l in scoped_leases if l.is_mutation]
|
mutations = [l for l in lease_impacts if l.is_mutation]
|
||||||
terminal_lock_in_scope = (
|
|
||||||
terminal_lock
|
|
||||||
if resolved_class
|
|
||||||
not in {
|
|
||||||
RestartClass.CLIENT_RECONNECT,
|
|
||||||
RestartClass.SESSION_RECONNECT,
|
|
||||||
}
|
|
||||||
else None
|
|
||||||
)
|
|
||||||
|
|
||||||
affected_issues = sorted(
|
affected_issues = sorted(
|
||||||
{
|
{
|
||||||
@@ -672,24 +348,9 @@ def evaluate_restart_impact(
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
disruptive = bool(
|
disruptive = bool(disruptive_leases or other_live_sessions or terminal_lock)
|
||||||
disruptive_leases or other_live_sessions or terminal_lock_in_scope
|
|
||||||
)
|
|
||||||
|
|
||||||
authorization_ok = bool(
|
if not inventory_complete:
|
||||||
not unknown_class
|
|
||||||
and permission_authorized
|
|
||||||
and role_authorized
|
|
||||||
and approval_satisfied
|
|
||||||
and target_complete
|
|
||||||
)
|
|
||||||
|
|
||||||
if policy_enforced and not authorization_ok:
|
|
||||||
verdict = VERDICT_UNSAFE
|
|
||||||
allow_restart = False
|
|
||||||
reasons.append("restart class authorization denied (fail closed)")
|
|
||||||
reasons.extend(authorization_reasons)
|
|
||||||
elif not inventory_complete:
|
|
||||||
verdict = VERDICT_UNSAFE
|
verdict = VERDICT_UNSAFE
|
||||||
allow_restart = False
|
allow_restart = False
|
||||||
reasons.append(
|
reasons.append(
|
||||||
@@ -720,7 +381,7 @@ def evaluate_restart_impact(
|
|||||||
f"{len(critical_sections)} critical section(s) in flight "
|
f"{len(critical_sections)} critical section(s) in flight "
|
||||||
"(active lease with a live owner)"
|
"(active lease with a live owner)"
|
||||||
)
|
)
|
||||||
if terminal_lock_in_scope:
|
if terminal_lock:
|
||||||
reasons.append("active terminal (merge) lock present")
|
reasons.append("active terminal (merge) lock present")
|
||||||
|
|
||||||
override_would_allow = bool(inventory_complete and disruptive)
|
override_would_allow = bool(inventory_complete and disruptive)
|
||||||
@@ -750,12 +411,6 @@ def evaluate_restart_impact(
|
|||||||
audit_record = {
|
audit_record = {
|
||||||
"event": "restart_impact_evaluated",
|
"event": "restart_impact_evaluated",
|
||||||
"coordinator_version": COORDINATOR_VERSION,
|
"coordinator_version": COORDINATOR_VERSION,
|
||||||
"restart_class": (
|
|
||||||
resolved_class.value if resolved_class else str(restart_class or "")
|
|
||||||
),
|
|
||||||
"required_permission": (
|
|
||||||
policy.required_permission if policy is not None else None
|
|
||||||
),
|
|
||||||
"evaluated_at": moment.isoformat(),
|
"evaluated_at": moment.isoformat(),
|
||||||
"dry_run": dry_run,
|
"dry_run": dry_run,
|
||||||
"operator_override": bool(operator_override),
|
"operator_override": bool(operator_override),
|
||||||
@@ -769,15 +424,6 @@ def evaluate_restart_impact(
|
|||||||
|
|
||||||
return RestartImpactReport(
|
return RestartImpactReport(
|
||||||
coordinator_version=COORDINATOR_VERSION,
|
coordinator_version=COORDINATOR_VERSION,
|
||||||
restart_class=(
|
|
||||||
resolved_class.value if resolved_class else str(restart_class or "")
|
|
||||||
),
|
|
||||||
restart_policy=policy.as_dict() if policy is not None else {},
|
|
||||||
policy_enforced=policy_enforced,
|
|
||||||
permission_authorized=permission_authorized,
|
|
||||||
role_authorized=role_authorized,
|
|
||||||
approval_satisfied=approval_satisfied,
|
|
||||||
authorization_reasons=authorization_reasons,
|
|
||||||
evaluated_at=moment.isoformat(),
|
evaluated_at=moment.isoformat(),
|
||||||
dry_run=dry_run,
|
dry_run=dry_run,
|
||||||
restart_performed=False,
|
restart_performed=False,
|
||||||
@@ -794,11 +440,9 @@ def evaluate_restart_impact(
|
|||||||
affected_issues=affected_issues,
|
affected_issues=affected_issues,
|
||||||
affected_prs=affected_prs,
|
affected_prs=affected_prs,
|
||||||
mutations=mutations,
|
mutations=mutations,
|
||||||
terminal_lock=(
|
terminal_lock=dict(terminal_lock)
|
||||||
dict(terminal_lock_in_scope)
|
if isinstance(terminal_lock, Mapping)
|
||||||
if isinstance(terminal_lock_in_scope, Mapping)
|
else terminal_lock,
|
||||||
else terminal_lock_in_scope
|
|
||||||
),
|
|
||||||
ack_state=ack_state,
|
ack_state=ack_state,
|
||||||
prior_recovery_attempts=prior_recovery_attempts,
|
prior_recovery_attempts=prior_recovery_attempts,
|
||||||
counts=counts,
|
counts=counts,
|
||||||
|
|||||||
@@ -0,0 +1,596 @@
|
|||||||
|
"""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)
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
@@ -1,232 +0,0 @@
|
|||||||
"""Permission, drain, routing, and audit matrix for restart classes (#663)."""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import os
|
|
||||||
from datetime import datetime, timezone
|
|
||||||
|
|
||||||
import restart_coordinator as rc
|
|
||||||
|
|
||||||
NOW = datetime(2026, 7, 24, 20, 0, tzinfo=timezone.utc)
|
|
||||||
|
|
||||||
|
|
||||||
def _inventory() -> dict:
|
|
||||||
return {
|
|
||||||
"inventory_complete": True,
|
|
||||||
"sessions": [
|
|
||||||
{
|
|
||||||
"session_id": "requester",
|
|
||||||
"role": "author",
|
|
||||||
"profile": "prgs-author",
|
|
||||||
"pid": os.getpid(),
|
|
||||||
"status": "active",
|
|
||||||
"last_heartbeat_at": NOW.isoformat(),
|
|
||||||
},
|
|
||||||
{
|
|
||||||
"session_id": "reviewer",
|
|
||||||
"role": "reviewer",
|
|
||||||
"profile": "prgs-reviewer",
|
|
||||||
"pid": os.getpid(),
|
|
||||||
"status": "active",
|
|
||||||
"last_heartbeat_at": NOW.isoformat(),
|
|
||||||
},
|
|
||||||
],
|
|
||||||
"leases": [
|
|
||||||
{
|
|
||||||
"lease_id": "review-lease",
|
|
||||||
"session_id": "reviewer",
|
|
||||||
"role": "reviewer",
|
|
||||||
"phase": "reviewing",
|
|
||||||
"work_kind": "pr",
|
|
||||||
"work_number": 900,
|
|
||||||
"worktree_path": "/tmp/review-900",
|
|
||||||
"freshness": {"freshness": "active"},
|
|
||||||
}
|
|
||||||
],
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
def _evaluate(
|
|
||||||
restart_class: rc.RestartClass,
|
|
||||||
*,
|
|
||||||
role: str = "controller",
|
|
||||||
permissions: tuple[str, ...] | None = None,
|
|
||||||
approved: bool = True,
|
|
||||||
operator: bool = True,
|
|
||||||
**targets,
|
|
||||||
):
|
|
||||||
return rc.evaluate_restart_impact(
|
|
||||||
_inventory(),
|
|
||||||
now=NOW,
|
|
||||||
requesting_session_id="requester",
|
|
||||||
restart_class=restart_class,
|
|
||||||
requester_role=role,
|
|
||||||
requester_permissions=(
|
|
||||||
permissions if permissions is not None
|
|
||||||
else rc.permissions_for_role(role)
|
|
||||||
),
|
|
||||||
controller_approved=approved,
|
|
||||||
operator_authorized=operator,
|
|
||||||
**targets,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def test_policy_table_covers_exactly_all_nine_classes():
|
|
||||||
assert set(rc.RESTART_CLASS_POLICIES) == set(rc.RestartClass)
|
|
||||||
assert len(rc.RESTART_CLASS_POLICIES) == 9
|
|
||||||
for restart_class, policy in rc.RESTART_CLASS_POLICIES.items():
|
|
||||||
assert policy.restart_class is restart_class
|
|
||||||
assert policy.required_permission
|
|
||||||
assert policy.expected_blast_radius in {
|
|
||||||
rc.BLAST_NONE, rc.BLAST_LOW, rc.BLAST_MEDIUM, rc.BLAST_HIGH
|
|
||||||
}
|
|
||||||
assert policy.drain_requirement
|
|
||||||
assert policy.approval_requirement
|
|
||||||
assert policy.audit_requirement
|
|
||||||
assert policy.recovery_behavior
|
|
||||||
|
|
||||||
|
|
||||||
def test_permission_matrix_allows_each_class_with_exact_permission():
|
|
||||||
targets = {
|
|
||||||
rc.RestartClass.WORKER_RESTART: {"target_session_id": "reviewer"},
|
|
||||||
rc.RestartClass.ROLE_RUNTIME_RESTART: {"target_role": "reviewer"},
|
|
||||||
rc.RestartClass.CONNECTOR_RESTART: {"target_connector": "github"},
|
|
||||||
}
|
|
||||||
for restart_class, policy in rc.RESTART_CLASS_POLICIES.items():
|
|
||||||
report = _evaluate(
|
|
||||||
restart_class,
|
|
||||||
permissions=(policy.required_permission,),
|
|
||||||
**targets.get(restart_class, {}),
|
|
||||||
)
|
|
||||||
assert report.permission_authorized, restart_class
|
|
||||||
assert report.role_authorized, restart_class
|
|
||||||
assert report.approval_satisfied, restart_class
|
|
||||||
assert report.audit_record["restart_class"] == restart_class.value
|
|
||||||
assert (
|
|
||||||
report.audit_record["required_permission"]
|
|
||||||
== policy.required_permission
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def test_missing_or_nearby_permission_denies():
|
|
||||||
report = _evaluate(
|
|
||||||
rc.RestartClass.ROLE_RUNTIME_RESTART,
|
|
||||||
permissions=("mcp.restart.worker.request",),
|
|
||||||
target_role="reviewer",
|
|
||||||
)
|
|
||||||
assert report.verdict == rc.VERDICT_UNSAFE
|
|
||||||
assert not report.allow_restart
|
|
||||||
assert not report.permission_authorized
|
|
||||||
assert any("missing required permission" in r for r in report.reasons)
|
|
||||||
|
|
||||||
|
|
||||||
def test_unknown_restart_class_denies_fail_closed():
|
|
||||||
report = rc.evaluate_restart_impact(
|
|
||||||
_inventory(),
|
|
||||||
now=NOW,
|
|
||||||
restart_class="surprise_reboot",
|
|
||||||
requester_role="admin",
|
|
||||||
requester_permissions=("mcp.restart.host.request",),
|
|
||||||
controller_approved=True,
|
|
||||||
operator_authorized=True,
|
|
||||||
)
|
|
||||||
assert report.verdict == rc.VERDICT_UNSAFE
|
|
||||||
assert not report.allow_restart
|
|
||||||
assert report.restart_policy == {}
|
|
||||||
assert any("unknown restart class" in r for r in report.reasons)
|
|
||||||
|
|
||||||
|
|
||||||
def test_worker_roles_cannot_request_full_or_host_restart():
|
|
||||||
for role in rc.WORKER_ROLES:
|
|
||||||
granted = rc.permissions_for_role(role)
|
|
||||||
assert "mcp.restart.full.request" not in granted
|
|
||||||
assert "mcp.restart.host.request" not in granted
|
|
||||||
report = _evaluate(
|
|
||||||
rc.RestartClass.FULL_MCP_RESTART,
|
|
||||||
role=role,
|
|
||||||
permissions=granted,
|
|
||||||
)
|
|
||||||
assert not report.role_authorized
|
|
||||||
assert not report.allow_restart
|
|
||||||
|
|
||||||
|
|
||||||
def test_controller_approval_is_independent_of_permission():
|
|
||||||
report = _evaluate(
|
|
||||||
rc.RestartClass.WORKER_RESTART,
|
|
||||||
approved=False,
|
|
||||||
target_session_id="reviewer",
|
|
||||||
)
|
|
||||||
assert report.permission_authorized
|
|
||||||
assert not report.approval_satisfied
|
|
||||||
assert not report.allow_restart
|
|
||||||
|
|
||||||
|
|
||||||
def test_narrow_classes_do_not_inherit_full_drain_or_peer_lease_block():
|
|
||||||
for restart_class in (
|
|
||||||
rc.RestartClass.CLIENT_RECONNECT,
|
|
||||||
rc.RestartClass.SESSION_RECONNECT,
|
|
||||||
rc.RestartClass.CONFIGURATION_RELOAD,
|
|
||||||
):
|
|
||||||
report = _evaluate(restart_class)
|
|
||||||
assert not report.restart_policy["full_drain_required"]
|
|
||||||
assert report.counts["leases_disruptive"] == 0
|
|
||||||
assert report.counts["sessions_live_other"] == 0
|
|
||||||
assert report.counts["critical_sections"] == 0
|
|
||||||
assert report.counts["mutations"] == 0
|
|
||||||
assert report.allow_restart, (restart_class, report.reasons)
|
|
||||||
|
|
||||||
|
|
||||||
def test_client_reconnect_does_not_wait_for_unrelated_terminal_lock():
|
|
||||||
inventory = _inventory()
|
|
||||||
inventory["terminal_lock"] = {"terminal_pr": 901}
|
|
||||||
report = rc.evaluate_restart_impact(
|
|
||||||
inventory,
|
|
||||||
now=NOW,
|
|
||||||
requesting_session_id="requester",
|
|
||||||
restart_class=rc.RestartClass.CLIENT_RECONNECT,
|
|
||||||
requester_role="author",
|
|
||||||
requester_permissions=rc.permissions_for_role("author"),
|
|
||||||
)
|
|
||||||
assert report.allow_restart
|
|
||||||
assert report.terminal_lock is None
|
|
||||||
|
|
||||||
|
|
||||||
def test_scoped_restart_only_counts_named_target():
|
|
||||||
report = _evaluate(
|
|
||||||
rc.RestartClass.ROLE_RUNTIME_RESTART,
|
|
||||||
target_role="author",
|
|
||||||
)
|
|
||||||
assert report.counts["leases_disruptive"] == 0
|
|
||||||
assert report.affected_prs == []
|
|
||||||
assert report.allow_restart
|
|
||||||
|
|
||||||
reviewer = _evaluate(
|
|
||||||
rc.RestartClass.ROLE_RUNTIME_RESTART,
|
|
||||||
target_role="reviewer",
|
|
||||||
)
|
|
||||||
assert reviewer.counts["leases_disruptive"] == 1
|
|
||||||
assert reviewer.affected_prs == [900]
|
|
||||||
assert not reviewer.allow_restart
|
|
||||||
|
|
||||||
|
|
||||||
def test_missing_scoped_target_denies_instead_of_widening():
|
|
||||||
for restart_class in (
|
|
||||||
rc.RestartClass.WORKER_RESTART,
|
|
||||||
rc.RestartClass.ROLE_RUNTIME_RESTART,
|
|
||||||
rc.RestartClass.CONNECTOR_RESTART,
|
|
||||||
):
|
|
||||||
report = _evaluate(restart_class)
|
|
||||||
assert not report.allow_restart
|
|
||||||
assert any("target required" in r for r in report.reasons)
|
|
||||||
|
|
||||||
|
|
||||||
def test_only_full_and_host_classes_require_full_drain():
|
|
||||||
requiring_full = {
|
|
||||||
restart_class
|
|
||||||
for restart_class, policy in rc.RESTART_CLASS_POLICIES.items()
|
|
||||||
if policy.full_drain_required
|
|
||||||
}
|
|
||||||
assert requiring_full == {
|
|
||||||
rc.RestartClass.FULL_MCP_RESTART,
|
|
||||||
rc.RestartClass.HOST_RESTART,
|
|
||||||
}
|
|
||||||
@@ -444,12 +444,6 @@ class TestAuditEmission(unittest.TestCase):
|
|||||||
)
|
)
|
||||||
self.assertEqual(record["target"]["namespace"], NAMESPACE)
|
self.assertEqual(record["target"]["namespace"], NAMESPACE)
|
||||||
self.assertEqual(record["target"]["mode"], "restart")
|
self.assertEqual(record["target"]["mode"], "restart")
|
||||||
self.assertEqual(
|
|
||||||
record["target"]["restart_class"], "role_runtime_restart"
|
|
||||||
)
|
|
||||||
self.assertEqual(
|
|
||||||
record["metadata"]["restart_class"], "role_runtime_restart"
|
|
||||||
)
|
|
||||||
self.assertEqual(record["result"], console_audit.RESULT_ALLOWED)
|
self.assertEqual(record["result"], console_audit.RESULT_ALLOWED)
|
||||||
self.assertEqual(record["actor"]["subject"], "[email protected]")
|
self.assertEqual(record["actor"]["subject"], "[email protected]")
|
||||||
self.assertFalse(record["metadata"]["process_kill_executed"])
|
self.assertFalse(record["metadata"]["process_kill_executed"])
|
||||||
|
|||||||
@@ -38,7 +38,6 @@ from dataclasses import asdict, dataclass
|
|||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
import mcp_namespace_health
|
import mcp_namespace_health
|
||||||
import restart_coordinator
|
|
||||||
import runtime_recovery_guard
|
import runtime_recovery_guard
|
||||||
from webui import console_audit, console_authz
|
from webui import console_audit, console_authz
|
||||||
|
|
||||||
@@ -100,14 +99,6 @@ def _clean(value: Any) -> str:
|
|||||||
return str(value or "").strip()
|
return str(value or "").strip()
|
||||||
|
|
||||||
|
|
||||||
def restart_class_for_mode(mode: str) -> str:
|
|
||||||
"""Map the existing namespace controls onto the #663 class taxonomy."""
|
|
||||||
|
|
||||||
if _clean(mode) == MODE_RELOAD:
|
|
||||||
return restart_coordinator.RestartClass.CONFIGURATION_RELOAD.value
|
|
||||||
return restart_coordinator.RestartClass.ROLE_RUNTIME_RESTART.value
|
|
||||||
|
|
||||||
|
|
||||||
# --- Mutation ledger --------------------------------------------------------
|
# --- Mutation ledger --------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
@@ -265,7 +256,6 @@ def build_restart_preview(
|
|||||||
|
|
||||||
return {
|
return {
|
||||||
"action_id": action_id,
|
"action_id": action_id,
|
||||||
"restart_class": restart_class_for_mode(md),
|
|
||||||
"namespace": ns,
|
"namespace": ns,
|
||||||
"mode": md,
|
"mode": md,
|
||||||
"scope_valid": scope_error is None,
|
"scope_valid": scope_error is None,
|
||||||
@@ -319,7 +309,6 @@ def assess_restart_request(
|
|||||||
"reason_code": reason_code,
|
"reason_code": reason_code,
|
||||||
"detail": detail,
|
"detail": detail,
|
||||||
"action_id": action_id,
|
"action_id": action_id,
|
||||||
"restart_class": restart_class_for_mode(md),
|
|
||||||
"namespace": ns,
|
"namespace": ns,
|
||||||
"mode": md,
|
"mode": md,
|
||||||
"preview": preview,
|
"preview": preview,
|
||||||
@@ -404,7 +393,6 @@ def assess_restart_request(
|
|||||||
"process."
|
"process."
|
||||||
),
|
),
|
||||||
"action_id": action_id,
|
"action_id": action_id,
|
||||||
"restart_class": restart_class_for_mode(md),
|
|
||||||
"namespace": ns,
|
"namespace": ns,
|
||||||
"mode": md,
|
"mode": md,
|
||||||
"preview": preview,
|
"preview": preview,
|
||||||
@@ -453,11 +441,7 @@ def execute_restart(
|
|||||||
else console_audit.RESULT_DENIED
|
else console_audit.RESULT_DENIED
|
||||||
),
|
),
|
||||||
principal=principal,
|
principal=principal,
|
||||||
target={
|
target={"namespace": assessment["namespace"], "mode": assessment["mode"]},
|
||||||
"namespace": assessment["namespace"],
|
|
||||||
"mode": assessment["mode"],
|
|
||||||
"restart_class": assessment["restart_class"],
|
|
||||||
},
|
|
||||||
reason_code=assessment["reason_code"],
|
reason_code=assessment["reason_code"],
|
||||||
detail=assessment["detail"],
|
detail=assessment["detail"],
|
||||||
request_id=request_id,
|
request_id=request_id,
|
||||||
@@ -466,7 +450,6 @@ def execute_restart(
|
|||||||
"gates_passed": assessment["gates_passed"],
|
"gates_passed": assessment["gates_passed"],
|
||||||
"process_kill_executed": False,
|
"process_kill_executed": False,
|
||||||
"post_restart_verification_required": True,
|
"post_restart_verification_required": True,
|
||||||
"restart_class": assessment["restart_class"],
|
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -480,7 +463,6 @@ def execute_restart(
|
|||||||
"namespace": assessment["namespace"],
|
"namespace": assessment["namespace"],
|
||||||
"mode": assessment["mode"],
|
"mode": assessment["mode"],
|
||||||
"action_id": action_id,
|
"action_id": action_id,
|
||||||
"restart_class": assessment["restart_class"],
|
|
||||||
"process_kill_executed": False,
|
"process_kill_executed": False,
|
||||||
"host_hook": assessment["preview"]["restart_hook"],
|
"host_hook": assessment["preview"]["restart_hook"],
|
||||||
"next_action": (
|
"next_action": (
|
||||||
|
|||||||
Reference in New Issue
Block a user