Files
Gitea-Tools/drain_proof.py
sysadmin 3a9d634c17 fix(drain-proof): bind acknowledgement coverage to session identity (#661)
Review 582 (REQUEST_CHANGES at 95178349) found a residual fail-open of the
same class the PR set out to close. Acknowledgement coverage was decided by
comparing a count against a count:

    covers_live_sessions = live_count_known and acked_count >= sessions_live_other

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

Coverage is now derived from authoritative impact-report evidence:

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

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

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

Reviewer's reproduction, before and after this commit:

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

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

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

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

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

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01VRUZAf3Fr5n3kqhhiayN6C
(cherry picked from commit 4193b63f415b066ee292386c2c89bc3d2651a0cc)
2026-07-25 00:04:17 -04:00

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,
)