Review 582 (REQUEST_CHANGES at95178349) found a residual fail-open of the same class the PR set out to close. Acknowledgement coverage was decided by comparing a count against a count: covers_live_sessions = live_count_known and acked_count >= sessions_live_other Nothing bound an acknowledgement to the identity of a session that actually owed one, so acknowledgements supplied for the requesting session and for a session that does not exist satisfied the obligations of two live sessions that never answered - minting a clean, correctly signed proof and an allow verdict from the restart gate. Coverage is now derived from authoritative impact-report evidence: - New `_required_ack_sessions()` derives the required session ids from the report itself, via `ack_state` keys and/or `affected_sessions` filtered on `live and not is_requester`. The requester is excluded only on explicit `is_requester` evidence, never inferred. - When both views are present they must name the same set, and the result is reconciled against `counts.sessions_live_other`. Missing, malformed, duplicated, contradictory, or unreconcilable identity evidence fails closed and outranks every permitting path, including the timeout policy. - Coverage requires every required id to carry an explicit acknowledgement token keyed by that id. Acknowledgements for the requester, for unknown ids, or for fabricated ids never increase coverage. - Caller-supplied acknowledgement cardinality is no longer proof of anything. Failure propagates unchanged through `acks_or_timeout` -> `proof.clean` -> `failed_checks` -> `gate_apply_restart` verdict `deny` / `allow=False`. The earlier missing-acknowledgement remediation is preserved in full: absent, None, non-mapping, empty, partial, stale, and unparseable acks still fail closed, `ack_timeout_policy_applied` stays strict `value is True`, and the legitimate zero-live-sessions and explicit-timeout paths still pass. Reviewer's reproduction, before and after this commit: sessions_live_other = 2 report ack_state = {'other-0': 'pending', 'other-1': 'pending'} supplied acks = {'req': 'ack', 'totally-bogus-session': 'ack'} before: acks_or_timeout = True | proof.clean = True | gate allow after: acks_or_timeout = False | proof.clean = False | gate deny Tests: 22 new cases in `AcknowledgementIdentityBindingTests` covering the reviewer's exact exploit, wrong-ids-with-sufficient-count, partial identity match, requester-only acks, fabricated ids, unproven per-session states, missing/malformed/contradictory identity evidence, count mismatch, and the preserved success paths. Verification: - `pytest tests/test_drain_proof.py` -> 61 passed, 56 subtests (baseline at95178349: 39 passed, 26 subtests) - Restart surface (6 modules) -> 154 passed, 68 subtests, exit 0 (baseline at95178349: 132 passed, 38 subtests) - Full `pytest tests/` -> 5177 passed vs baseline 5155 passed; the 23 failures are identical in both runs and pre-exist at95178349. Scope: drain_proof.py, tests/test_drain_proof.py. Co-Authored-By: Claude Opus 5 (1M context) <[email protected]> Claude-Session: https://claude.ai/code/session_01VRUZAf3Fr5n3kqhhiayN6C (cherry picked from commit 4193b63f415b066ee292386c2c89bc3d2651a0cc)
1016 lines
37 KiB
Python
1016 lines
37 KiB
Python
"""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 _session_ids_from_affected_sessions(
|
|
raw: Any,
|
|
) -> tuple[set[str] | None, str | None]:
|
|
"""Live, non-requester session ids from the report's ``affected_sessions``.
|
|
|
|
The requester is excluded only on authoritative ``is_requester`` evidence;
|
|
an entry whose ``live`` or ``is_requester`` flag is absent or not a real
|
|
bool is malformed, never assumed. Returns ``(ids, None)`` or
|
|
``(None, detail)``.
|
|
"""
|
|
|
|
if not isinstance(raw, (list, tuple)):
|
|
return None, (
|
|
"impact report 'affected_sessions' is malformed: expected a list, "
|
|
f"got {type(raw).__name__} (fail closed)"
|
|
)
|
|
collected: list[str] = []
|
|
for index, entry in enumerate(raw):
|
|
if not isinstance(entry, Mapping):
|
|
return None, (
|
|
f"impact report 'affected_sessions[{index}]' is malformed: "
|
|
f"expected a mapping, got {type(entry).__name__} (fail closed)"
|
|
)
|
|
session_id = entry.get("session_id")
|
|
if not isinstance(session_id, str) or not session_id.strip():
|
|
return None, (
|
|
f"impact report 'affected_sessions[{index}]' has no usable "
|
|
"session_id (fail closed)"
|
|
)
|
|
live = entry.get("live")
|
|
is_requester = entry.get("is_requester")
|
|
if not isinstance(live, bool) or not isinstance(is_requester, bool):
|
|
return None, (
|
|
f"impact report 'affected_sessions[{index}]' "
|
|
f"({session_id.strip()}) does not prove live/is_requester with "
|
|
"explicit booleans (fail closed)"
|
|
)
|
|
if live and not is_requester:
|
|
collected.append(session_id.strip())
|
|
if len(set(collected)) != len(collected):
|
|
return None, (
|
|
"impact report 'affected_sessions' names the same live session more "
|
|
"than once; required acknowledgement identities are ambiguous "
|
|
"(fail closed)"
|
|
)
|
|
return set(collected), None
|
|
|
|
|
|
def _session_ids_from_ack_state(raw: Any) -> tuple[set[str] | None, str | None]:
|
|
"""Session ids keyed by the report's ``ack_state``.
|
|
|
|
``ack_state`` is minted as ``{s.session_id: "pending" for s in
|
|
other_live_sessions}``, so its keys *are* the required set. Only the keys
|
|
are trusted here; the per-session value is the report's own placeholder and
|
|
is never read as an acknowledgement (acknowledgements come from the drain
|
|
state and must be explicit).
|
|
"""
|
|
|
|
if not isinstance(raw, Mapping):
|
|
return None, (
|
|
"impact report 'ack_state' is malformed: expected a mapping, got "
|
|
f"{type(raw).__name__} (fail closed)"
|
|
)
|
|
keys: set[str] = set()
|
|
for key in raw:
|
|
if not isinstance(key, str) or not key.strip():
|
|
return None, (
|
|
"impact report 'ack_state' carries a non-string or empty "
|
|
"session id (fail closed)"
|
|
)
|
|
keys.add(key.strip())
|
|
if len(keys) != len(raw):
|
|
return None, (
|
|
"impact report 'ack_state' names the same live session more than "
|
|
"once; required acknowledgement identities are ambiguous "
|
|
"(fail closed)"
|
|
)
|
|
return keys, None
|
|
|
|
|
|
def _required_ack_sessions(
|
|
impact_report: Mapping[str, Any],
|
|
) -> tuple[frozenset[str] | None, str | None]:
|
|
"""Identities of the live, non-requester sessions that owe an acknowledgement.
|
|
|
|
Derived only from authoritative impact-report evidence, never from the
|
|
caller-supplied drain state — the same principle checks 1 and 5 already
|
|
apply. Returns ``(ids, None)`` when the report proves the required set, or
|
|
``(None, detail)`` when the evidence is missing, malformed, contradictory,
|
|
or cannot be reconciled; every one of those fails the checklist closed.
|
|
|
|
The report carries two independent views of the same set and both are
|
|
validated: ``affected_sessions`` filtered on ``live and not is_requester``,
|
|
and ``ack_state`` whose keys are exactly those sessions. When both are
|
|
present they must name the same set, and the result is reconciled against
|
|
``counts.sessions_live_other`` — a report that disagrees with itself can
|
|
never authorise a restart.
|
|
"""
|
|
|
|
from_sessions: set[str] | None = None
|
|
if "affected_sessions" in impact_report:
|
|
from_sessions, detail = _session_ids_from_affected_sessions(
|
|
impact_report.get("affected_sessions")
|
|
)
|
|
if detail is not None:
|
|
return None, detail
|
|
|
|
from_ack_state: set[str] | None = None
|
|
if "ack_state" in impact_report:
|
|
from_ack_state, detail = _session_ids_from_ack_state(
|
|
impact_report.get("ack_state")
|
|
)
|
|
if detail is not None:
|
|
return None, detail
|
|
|
|
if from_sessions is None and from_ack_state is None:
|
|
return None, (
|
|
"impact report carries no session-identity evidence (neither "
|
|
"ack_state nor affected_sessions); required acknowledgements "
|
|
"cannot be attributed to a session (fail closed)"
|
|
)
|
|
|
|
if (
|
|
from_sessions is not None
|
|
and from_ack_state is not None
|
|
and from_sessions != from_ack_state
|
|
):
|
|
only_ack_state = sorted(from_ack_state - from_sessions)
|
|
only_sessions = sorted(from_sessions - from_ack_state)
|
|
return None, (
|
|
"impact report contradicts itself about which sessions must "
|
|
f"acknowledge: ack_state-only={only_ack_state}, "
|
|
f"affected_sessions-only={only_sessions} (fail closed)"
|
|
)
|
|
|
|
required = from_ack_state if from_ack_state is not None else from_sessions
|
|
if required is None: # unreachable; guarded above, kept fail-closed
|
|
return None, "required acknowledgement identities unresolved (fail closed)"
|
|
|
|
live_count = _live_session_count(impact_report)
|
|
if live_count is None:
|
|
return None, (
|
|
"impact report does not prove the live-session count; required "
|
|
"acknowledgement identities cannot be reconciled (fail closed)"
|
|
)
|
|
if len(required) != live_count:
|
|
return None, (
|
|
f"impact report names {len(required)} live session(s) requiring "
|
|
f"acknowledgement but counts.sessions_live_other={live_count}; "
|
|
"identity evidence cannot be reconciled (fail closed)"
|
|
)
|
|
return frozenset(required), None
|
|
|
|
|
|
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*, and *which sessions owe them*,
|
|
# are both answered by the impact report — never inferred from the shape or
|
|
# the cardinality of the drain state. Two fail-opens of the same class were
|
|
# closed here in turn. First, an absent ``acks`` key collapsed to ``{}`` and
|
|
# was read as "no session needed to acknowledge", so a proof minted clean
|
|
# while the report still showed other live sessions. Then coverage was
|
|
# decided by comparing an acknowledgement *count* against
|
|
# ``sessions_live_other``, so acknowledgements from the requesting session
|
|
# or from ids that do not exist satisfied the obligations of the live
|
|
# sessions that never answered.
|
|
#
|
|
# Coverage is therefore bound to identity: every live, non-requester session
|
|
# the report names must itself carry an explicit acknowledgement. An
|
|
# acknowledgement from any other id — the requester, an unknown id, a
|
|
# fabricated one — is not evidence about a required session and never
|
|
# increases coverage.
|
|
sessions_live_other = _live_session_count(impact_report)
|
|
required_sessions, identity_detail = _required_ack_sessions(impact_report)
|
|
# An unknown count or unresolvable identity evidence fails closed: neither
|
|
# can prove nobody had to acknowledge, so acknowledgement stays required.
|
|
acks_required = required_sessions is None or bool(required_sessions)
|
|
|
|
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)
|
|
# Identity-bound coverage. Only an entry keyed by a required session id and
|
|
# carrying an explicit token counts, so neither a requester acknowledgement
|
|
# nor an unknown id can stand in for a session that stayed silent.
|
|
if acks_present and required_sessions is not None:
|
|
acked_sessions = {
|
|
session_id
|
|
for session_id in required_sessions
|
|
if any(
|
|
isinstance(key, str)
|
|
and key.strip() == session_id
|
|
and _is_acknowledged(value)
|
|
for key, value in raw_acks.items()
|
|
)
|
|
}
|
|
else:
|
|
acked_sessions = set()
|
|
missing_sessions = (
|
|
sorted(required_sessions - acked_sessions)
|
|
if required_sessions is not None
|
|
else []
|
|
)
|
|
covers_live_sessions = required_sessions is not None and not missing_sessions
|
|
# 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 identity_detail is not None:
|
|
# Unresolvable identity evidence outranks every permitting path,
|
|
# including the timeout policy: when the report cannot say who owed an
|
|
# acknowledgement, nothing can show the obligation was discharged.
|
|
acks_ok = False
|
|
detail = identity_detail
|
|
elif 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 covers_live_sessions:
|
|
acks_ok = True
|
|
detail = (
|
|
f"{len(acked_sessions)} required live session(s) each acknowledged "
|
|
f"by identity: {sorted(acked_sessions)}; covers "
|
|
f"sessions_live_other={sessions_live_other}"
|
|
)
|
|
else:
|
|
acks_ok = False
|
|
if 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 = (
|
|
"no valid acknowledgement attributed to required live "
|
|
f"session(s) {missing_sessions}; the report shows "
|
|
f"sessions_live_other={sessions_live_other} and supplied "
|
|
"acknowledgements for other ids do not count (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,
|
|
)
|