#!/usr/bin/env python3 """Attended audit-core database lease failure/recovery exercise. The script is fail-closed and value-safe. It selects only a single live lease under the exact audit-core runtime prefix, never emits its id, never reads a database password, never restarts audit-core, and requires a separately approved synthetic-load driver. """ from __future__ import annotations import argparse import hashlib import importlib.util import json import os import stat import subprocess import sys import time from datetime import UTC, datetime, timedelta from pathlib import Path from typing import Any, Callable ROOT = Path(__file__).resolve().parents[1] TASK_ID = "RAILIANCE-WP-0024-T02" PROCEDURE = "audit-core-database-lease-recovery" LEASE_PREFIX = "database/creds/audit-core-runtime" EXTERNAL_SECRET = "audit-core-database" CONFIRM = f"{TASK_ID}:attended" MAX_WINDOW_SECONDS = 15 * 60 MIN_LEASE_TTL_SECONDS = 5 * 60 EXPECTED_DRIVER_KEYS = { "baseline": {"contract_id", "fixture_id", "status", "secret_values_observed"}, "expect-unavailable": { "contract_id", "fixture_id", "status", "http_status", "attempts", "secret_values_observed", }, "expect-recovered": { "contract_id", "fixture_id", "status", "http_status", "attempts", "secret_values_observed", }, "cleanup": {"contract_id", "fixture_id", "status", "secret_values_observed"}, } class ProcedureError(RuntimeError): pass def load_preflight_module() -> Any: spec = importlib.util.spec_from_file_location( "audit_core_recovery_preflight", ROOT / "scripts" / "audit-core-recovery-preflight.py", ) if not spec or not spec.loader: raise ProcedureError("cannot load recovery preflight module") module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) return module PREFLIGHT = load_preflight_module() def parse_time(value: str) -> datetime: parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) if parsed.tzinfo is None: raise ProcedureError("approval timestamps must include a timezone") return parsed.astimezone(UTC) def load_approval(path: Path, now: datetime, *, require_open_window: bool) -> dict[str, Any]: try: document = json.loads(path.read_text(encoding="utf-8")) except (OSError, json.JSONDecodeError) as exc: raise ProcedureError("approval receipt is unavailable or invalid JSON") from exc if not isinstance(document, dict): raise ProcedureError("approval receipt must be an object") if document.get("procedure") != PROCEDURE or document.get("task_id") != TASK_ID: raise ProcedureError("approval receipt has the wrong procedure or task") if document.get("status") != "approved" or not document.get("approval_id"): raise ProcedureError("approval receipt is not approved") window = document.get("window") or {} if not window.get("start") or not window.get("end"): raise ProcedureError("approval receipt has no complete window") start, end = parse_time(window["start"]), parse_time(window["end"]) if not start < end or (end - start).total_seconds() > MAX_WINDOW_SECONDS: raise ProcedureError("approval window must be positive and at most 15 minutes") if require_open_window and not start <= now <= end: raise ProcedureError("current time is outside the approved window") if not document.get("abort_operator"): raise ProcedureError("approval receipt has no abort operator") owners = document.get("owners") or {} for owner in ("audit-core", "rapp-postgres", "railiance-platform"): record = owners.get(owner) or {} if record.get("acknowledged") is not True or not record.get("message_id"): raise ProcedureError(f"approval receipt lacks {owner} acknowledgement") load = document.get("synthetic_load") or {} if not load.get("contract_id") or not load.get("driver_revision"): raise ProcedureError("approval receipt lacks the synthetic-load contract") return document def safe_run( command: list[str], *, label: str, env: dict[str, str] | None = None ) -> subprocess.CompletedProcess[str]: completed = subprocess.run( command, text=True, capture_output=True, env=env, check=False ) if completed.returncode != 0: # Never attach stdout/stderr: Bao and load-driver processes may hold # sensitive material even though their contract forbids emitting it. raise ProcedureError(f"{label} failed (exit {completed.returncode})") return completed class Authority: def __init__(self, token_file: Path) -> None: if not token_file.is_file() or stat.S_IMODE(token_file.stat().st_mode) != 0o600: raise ProcedureError("OpenBao token file must exist with mode 0600") token = token_file.read_text(encoding="utf-8").splitlines()[0].strip() if not token: raise ProcedureError("OpenBao token file is empty") self.env = dict(os.environ, BAO_ADDR="https://bao.coulomb.social", BAO_TOKEN=token) def bao(self, args: list[str], *, label: str) -> str: return safe_run(["bao", *args], label=label, env=self.env).stdout.strip() def single_runtime_lease(self) -> dict[str, Any]: raw = self.bao( ["list", "-format=json", f"sys/leases/lookup/{LEASE_PREFIX}"], label="list exact runtime lease handles", ) try: handles = json.loads(raw) except json.JSONDecodeError as exc: raise ProcedureError("runtime lease list is invalid JSON") from exc if not isinstance(handles, list) or len(handles) != 1: count = len(handles) if isinstance(handles, list) else "unknown" raise ProcedureError( f"exact runtime prefix must contain one live handle (observed {count})" ) lease_id = f"{LEASE_PREFIX}/{handles[0]}" lookup_raw = self.bao( ["write", "-format=json", "sys/leases/lookup", f"lease_id={lease_id}"], label="lookup exact runtime lease metadata", ) try: lookup = json.loads(lookup_raw)["data"] issue = parse_time(lookup["issue_time"]) expires = parse_time(lookup["expire_time"]) ttl = int(lookup["ttl"]) except (json.JSONDecodeError, KeyError, TypeError, ValueError) as exc: raise ProcedureError("runtime lease metadata is incomplete") from exc return { "id": lease_id, "fingerprint": hashlib.sha256(lease_id.encode()).hexdigest()[:12], "issue_time": issue, "expire_time": expires, "ttl": ttl, } def revoke(self, lease_id: str) -> None: self.bao(["lease", "revoke", lease_id], label="revoke exact runtime lease") class LoadDriver: def __init__(self, path: Path, contract_id: str, expected_revision: str) -> None: if not path.is_file() or not os.access(path, os.X_OK): raise ProcedureError("synthetic-load driver must be an executable file") revision = "sha256:" + hashlib.sha256(path.read_bytes()).hexdigest() if expected_revision != revision: raise ProcedureError("synthetic-load driver does not match the approved revision") self.path = path self.contract_id = contract_id self.revision = revision def run(self, phase: str) -> dict[str, Any]: result = safe_run( [str(self.path), phase, "--contract-id", self.contract_id], label=f"synthetic-load {phase}", ) try: payload = json.loads(result.stdout) except json.JSONDecodeError as exc: raise ProcedureError(f"synthetic-load {phase} returned invalid JSON") from exc if not isinstance(payload, dict) or set(payload) != EXPECTED_DRIVER_KEYS[phase]: raise ProcedureError(f"synthetic-load {phase} returned an unsafe evidence shape") if payload.get("contract_id") != self.contract_id: raise ProcedureError(f"synthetic-load {phase} returned the wrong contract") if payload.get("secret_values_observed") is not False: raise ProcedureError(f"synthetic-load {phase} did not attest value safety") return payload def secret_state(remote: Any) -> dict[str, Any]: secret_rv = remote.kubectl( ["-n", "audit-core", "get", "secret", EXTERNAL_SECRET, "-o", "jsonpath={.metadata.resourceVersion}"], label="read database Secret resource version", ) external = remote.kubectl_json( ["-n", "audit-core", "get", "externalsecret", EXTERNAL_SECRET], label="read database ExternalSecret state", ) mount_generation = remote.kubectl( ["-n", "audit-core", "exec", "deploy/audit-core", "--", "readlink", "/etc/audit-core/db/..data"], label="read mounted credential generation", ) refresh = external.get("status", {}).get("refreshTime") ready = PREFLIGHT.resource_condition(external)["ready"] if not refresh: raise ProcedureError("database ExternalSecret has no refresh time") return { "resource_version": secret_rv, "mount_generation": mount_generation, "refresh_time": parse_time(refresh), "ready": ready, } def pod_state(remote: Any) -> dict[str, Any]: pods = remote.kubectl_json( ["-n", "audit-core", "get", "pods", "-l", "app.kubernetes.io/name=audit-core"], label="read audit-core pod metadata", ).get("items", []) if len(pods) != 1: raise ProcedureError(f"expected one audit-core pod, observed {len(pods)}") statuses = pods[0].get("status", {}).get("containerStatuses", []) return { "uid": pods[0]["metadata"]["uid"], "restart_count": sum(int(item.get("restartCount", 0)) for item in statuses), } def assert_lease_matches_refresh(lease: dict[str, Any], secret: dict[str, Any]) -> None: delta = abs((lease["issue_time"] - secret["refresh_time"]).total_seconds()) if ( delta > 5 or secret["refresh_time"] < lease["issue_time"] - timedelta(seconds=5) or secret["refresh_time"] > lease["expire_time"] ): raise ProcedureError("single runtime lease is not coherent with the mounted Secret refresh") if lease["ttl"] < MIN_LEASE_TTL_SECONDS: raise ProcedureError("runtime lease is too close to automatic refresh; wait for the next sync") def wait_for( predicate: Callable[[], bool], *, label: str, timeout: int = 120, interval: float = 2 ) -> None: deadline = time.monotonic() + timeout while time.monotonic() < deadline: if predicate(): return time.sleep(interval) raise ProcedureError(f"timed out waiting for {label}") def force_refresh(remote: Any, baseline: dict[str, Any]) -> dict[str, Any]: remote.kubectl( [ "-n", "audit-core", "annotate", "externalsecret", EXTERNAL_SECRET, f"railiance.io/force-sync={int(time.time())}", "--overwrite", ], label="force database ExternalSecret reconciliation", ) latest: dict[str, Any] = {} def changed() -> bool: nonlocal latest latest = secret_state(remote) return bool( latest["ready"] and latest["resource_version"] != baseline["resource_version"] and latest["mount_generation"] != baseline["mount_generation"] and latest["refresh_time"] > baseline["refresh_time"] ) wait_for(changed, label="database Secret and mounted generation refresh") return latest def exercise(args: argparse.Namespace) -> dict[str, Any]: approval = load_approval(args.approval, datetime.now(UTC), require_open_window=True) if args.confirm != CONFIRM: raise ProcedureError(f"live exercise requires --confirm {CONFIRM}") remote = PREFLIGHT.Remote(args.remote) authority = Authority(args.token_file) driver = LoadDriver( args.load_driver, approval["synthetic_load"]["contract_id"], approval["synthetic_load"]["driver_revision"], ) preflight_args = argparse.Namespace( approved_window_id=approval["approval_id"], audit_core_owner_ack=True, rapp_postgres_owner_ack=True, synthetic_load_id=approval["synthetic_load"]["contract_id"], abort_operator=approval["abort_operator"], ) preflight = PREFLIGHT.database_lease_preflight(remote, preflight_args) if not preflight["ready_for_live_execution"]: raise ProcedureError("database lease recovery preflight is not ready") before_pod = pod_state(remote) before_secret = secret_state(remote) lease = authority.single_runtime_lease() assert_lease_matches_refresh(lease, before_secret) baseline_load = driver.run("baseline") if baseline_load.get("status") != "ready": raise ProcedureError("synthetic-load baseline is not ready") revoked = False recovered_secret: dict[str, Any] | None = None unavailable: dict[str, Any] | None = None recovered: dict[str, Any] | None = None cleanup_result: dict[str, Any] | None = None completed = False try: # Close the race with ESO's ordinary refresh before the destructive step. if secret_state(remote) != before_secret or pod_state(remote) != before_pod: raise ProcedureError("baseline changed before revocation") current = authority.single_runtime_lease() if current["fingerprint"] != lease["fingerprint"]: raise ProcedureError("runtime lease changed before revocation") authority.revoke(lease["id"]) revoked = True wait_for( lambda: PREFLIGHT.endpoint_status(remote, "/healthz") == 200 and PREFLIGHT.endpoint_status(remote, "/readyz") == 503, label="health 200 and readiness 503 after revocation", timeout=60, ) unavailable = driver.run("expect-unavailable") if ( unavailable.get("http_status") != 503 or unavailable.get("status") != "retryable_unavailable" or not isinstance(unavailable.get("attempts"), int) or unavailable["attempts"] < 1 or unavailable.get("fixture_id") != baseline_load.get("fixture_id") ): raise ProcedureError("synthetic load did not prove retryable 503") recovered_secret = force_refresh(remote, before_secret) wait_for( lambda: PREFLIGHT.endpoint_status(remote, "/readyz") == 200, label="audit-core readiness recovery", ) recovered = driver.run("expect-recovered") if ( recovered.get("status") not in {"accepted", "duplicate"} or recovered.get("http_status") not in {200, 202} or recovered.get("fixture_id") != baseline_load.get("fixture_id") ): raise ProcedureError("synthetic load did not prove accepted/duplicate recovery") after_pod = pod_state(remote) if after_pod != before_pod: raise ProcedureError("audit-core pod identity or restart count changed") completed = True finally: if revoked and recovered_secret is None: try: force_refresh(remote, before_secret) except Exception: pass try: cleanup_result = driver.run("cleanup") if ( cleanup_result.get("status") != "clean" or cleanup_result.get("fixture_id") != baseline_load.get("fixture_id") ): raise ProcedureError("synthetic-load cleanup did not confirm exact fixture removal") except Exception: if completed: raise return { "procedure": PROCEDURE, "task_id": TASK_ID, "approval_id": approval["approval_id"], "synthetic_load_contract": approval["synthetic_load"]["contract_id"], "synthetic_load_driver_revision": driver.revision, "lease_handle_fingerprint": lease["fingerprint"], "lease_revoked": revoked, "health_during_failure": 200, "readiness_during_failure": 503, "synthetic_unavailable_attempts": unavailable["attempts"] if unavailable else None, "secret_resource_version_changed": bool( recovered_secret and recovered_secret["resource_version"] != before_secret["resource_version"] ), "mount_generation_changed": bool( recovered_secret and recovered_secret["mount_generation"] != before_secret["mount_generation"] ), "readiness_recovered": PREFLIGHT.endpoint_status(remote, "/readyz") == 200, "synthetic_recovery_status": recovered["status"] if recovered else None, "same_pod_uid": pod_state(remote)["uid"] == before_pod["uid"], "restart_count_unchanged": pod_state(remote)["restart_count"] == before_pod["restart_count"], "load_cleanup_status": cleanup_result["status"] if cleanup_result else None, "completed_at": datetime.now(UTC).replace(microsecond=0).isoformat().replace("+00:00", "Z"), "secret_values_observed": False, } def main() -> int: parser = argparse.ArgumentParser() parser.add_argument("command", choices=["validate-approval", "exercise"]) parser.add_argument("--approval", type=Path, required=True) parser.add_argument("--remote", default="railiance01") parser.add_argument( "--token-file", type=Path, default=Path.home() / ".local/openbao/platform-admin.token", ) parser.add_argument("--load-driver", type=Path) parser.add_argument("--confirm") args = parser.parse_args() try: if args.command == "validate-approval": approval = load_approval(args.approval, datetime.now(UTC), require_open_window=False) result = { "procedure": PROCEDURE, "approval_id": approval["approval_id"], "approval_receipt_valid": True, "secret_values_observed": False, } else: if args.load_driver is None: raise ProcedureError("live exercise requires --load-driver") result = exercise(args) except (OSError, IndexError, ProcedureError, PREFLIGHT.PreflightError) as exc: print(f"database lease recovery failed: {exc}", file=sys.stderr) return 1 print(json.dumps(result, indent=2, sort_keys=True)) return 0 if __name__ == "__main__": raise SystemExit(main())