"""SQLite-backed approval object, closed machine, local outbox. All mutations run in BEGIN IMMEDIATE and insert the outbox row before COMMIT. Consumption follows GH-DEC-2026-003: the PEP presents the decision binding's request digest before the protected side effect. """ from __future__ import annotations import json import sqlite3 import threading from dataclasses import dataclass from datetime import datetime, timedelta, timezone from pathlib import Path from typing import Any, Callable from uuid import uuid4 from .binding import binding_digest, canonical_binding, require_digest from .errors import ( Conflict, DuplicateApprover, NotFound, StoreUnavailable, Unprocessable, ) ISSUER = "approval-engine" CLAIM_SCHEMA = "0.1" YIELDS_TO = "net-kingdom taxonomy request-claim schema (statute ยง17; unassigned)" DEFAULT_FRESHNESS_TTL = 30 AUDIT_SCHEMA = "audit-core.event.v1alpha1" SOURCE = "approval-engine" SCOPE = "netkingdom-approvals" EVENT_CLASSES = ("issuance", "use", "supersession", "revocation", "heartbeat") SCHEMA = """ CREATE TABLE IF NOT EXISTS approvals ( id TEXT PRIMARY KEY, status TEXT NOT NULL, binding_json TEXT NOT NULL, binding_digest TEXT NOT NULL, pdp_digest TEXT, actor TEXT NOT NULL, principal TEXT NOT NULL, action TEXT NOT NULL, purpose TEXT NOT NULL, target_json TEXT NOT NULL, not_before TEXT NOT NULL, expires_at TEXT NOT NULL, required_count INTEGER NOT NULL, superseded_by TEXT, consumed_digest TEXT, consumed_decision_id TEXT, consumed_at TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS entries ( approval_id TEXT NOT NULL, subject_id TEXT NOT NULL, approved_at TEXT NOT NULL, assurance TEXT, evidence_ref TEXT, PRIMARY KEY (approval_id, subject_id) ); CREATE TABLE IF NOT EXISTS outbox ( event_id TEXT PRIMARY KEY, class TEXT NOT NULL, approval_id TEXT, payload_json TEXT NOT NULL, created_at TEXT NOT NULL, drained_at TEXT ); """ ACTIVE = ("requested", "approved") MUTABLE = ("requested", "approved") def utc_now(clock: Callable[[], datetime] | None = None) -> datetime: if clock is not None: now = clock() else: now = datetime.now(timezone.utc) if now.tzinfo is None: now = now.replace(tzinfo=timezone.utc) return now.astimezone(timezone.utc).replace(microsecond=0) def iso(ts: datetime) -> str: return ts.astimezone(timezone.utc).replace(microsecond=0).isoformat() def parse_iso(value: str) -> datetime: return datetime.fromisoformat(value.replace("Z", "+00:00")) @dataclass class Approval: id: str status: str binding: dict[str, Any] binding_digest: str pdp_digest: str | None actor: str principal: str action: str purpose: str target: dict[str, Any] not_before: str expires_at: str required_count: int superseded_by: str | None consumed_digest: str | None consumed_decision_id: str | None consumed_at: str | None created_at: str updated_at: str entries: list[dict[str, Any]] def as_dict(self) -> dict[str, Any]: result = { "id": self.id, "status": self.status, "binding": { **self.binding, "digest": self.binding_digest, **({"pdp_digest": self.pdp_digest} if self.pdp_digest else {}), }, "validity": {"not_before": self.not_before, "expires_at": self.expires_at}, "required_count": self.required_count, "entries": self.entries, "superseded_by": self.superseded_by, "created_at": self.created_at, "updated_at": self.updated_at, } if self.consumed_digest: result["consumption"] = { "request_digest": self.consumed_digest, **( {"decision_id": self.consumed_decision_id} if self.consumed_decision_id else {} ), "consumed_at": self.consumed_at, } return result class Engine: def __init__( self, path: str | Path = ":memory:", *, clock: Callable[[], datetime] | None = None, freshness_ttl: int = DEFAULT_FRESHNESS_TTL, fail_outbox: bool = False, tenant: str = "platform", ) -> None: self.path = str(path) self.clock = clock self.freshness_ttl = freshness_ttl self.fail_outbox = fail_outbox self.tenant = tenant self._local = threading.local() self._init_schema() def _connect(self) -> sqlite3.Connection: try: conn = sqlite3.connect(self.path, timeout=5, check_same_thread=False) conn.row_factory = sqlite3.Row conn.execute("PRAGMA foreign_keys=ON") conn.execute("PRAGMA busy_timeout=5000") if self.path != ":memory:": conn.execute("PRAGMA journal_mode=WAL") return conn except sqlite3.Error as exc: raise StoreUnavailable(str(exc)) from exc def _conn(self) -> sqlite3.Connection: conn = getattr(self._local, "conn", None) if conn is None: conn = self._connect() self._local.conn = conn return conn def _init_schema(self) -> None: conn = self._conn() try: conn.executescript(SCHEMA) columns = { row["name"] for row in conn.execute("PRAGMA table_info(approvals)").fetchall() } for name in ("consumed_digest", "consumed_decision_id", "consumed_at"): if name not in columns: conn.execute(f"ALTER TABLE approvals ADD COLUMN {name} TEXT") conn.commit() except sqlite3.Error as exc: raise StoreUnavailable(str(exc)) from exc def close(self) -> None: conn = getattr(self._local, "conn", None) if conn is not None: conn.close() self._local.conn = None def now(self) -> datetime: return utc_now(self.clock) def create( self, binding: dict[str, Any], validity: dict[str, str], required_count: int = 1, *, pdp_digest: str | None = None, approval_id: str | None = None, ) -> Approval: if required_count < 1: raise Unprocessable("required_count must be >= 1") canon = canonical_binding(binding) digest = binding_digest(canon) pdp = require_digest(pdp_digest) not_before = validity.get("not_before") or iso(self.now()) expires_at = validity.get("expires_at") if not expires_at: raise Unprocessable("validity.expires_at is required") if parse_iso(expires_at) <= parse_iso(not_before): raise Unprocessable("expires_at must be after not_before") aid = approval_id or str(uuid4()) now = iso(self.now()) conn = self._conn() try: conn.execute("BEGIN IMMEDIATE") conn.execute( """INSERT INTO approvals ( id, status, binding_json, binding_digest, pdp_digest, actor, principal, action, purpose, target_json, not_before, expires_at, required_count, superseded_by, created_at, updated_at ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", ( aid, "requested", json.dumps(canon, sort_keys=True), digest, pdp, canon["actor"], canon["principal"], canon["action"], canon["purpose"], json.dumps(canon["target"], sort_keys=True), not_before, expires_at, required_count, None, now, now, ), ) conn.commit() except sqlite3.IntegrityError as exc: conn.rollback() raise Conflict("approval id already exists") from exc except sqlite3.Error as exc: conn.rollback() raise StoreUnavailable(str(exc)) from exc except Exception: conn.rollback() raise return self.get(aid) def get(self, approval_id: str) -> Approval: conn = self._conn() try: row = conn.execute("SELECT * FROM approvals WHERE id=?", (approval_id,)).fetchone() except sqlite3.Error as exc: raise StoreUnavailable(str(exc)) from exc if row is None: raise NotFound(approval_id) return self._hydrate(row, persist_expiry=True) def _hydrate(self, row: sqlite3.Row, persist_expiry: bool) -> Approval: entries = [ { "subject_id": e["subject_id"], "approved_at": e["approved_at"], "assurance": e["assurance"], "evidence_ref": e["evidence_ref"], } for e in self._conn() .execute( "SELECT subject_id, approved_at, assurance, evidence_ref " "FROM entries WHERE approval_id=? ORDER BY approved_at", (row["id"],), ) .fetchall() ] status = row["status"] if persist_expiry and status in ACTIVE: now = self.now() if now >= parse_iso(row["expires_at"]): status = self._expire(row["id"], status) binding = json.loads(row["binding_json"]) return Approval( id=row["id"], status=status, binding=binding, binding_digest=row["binding_digest"], pdp_digest=row["pdp_digest"], actor=row["actor"], principal=row["principal"], action=row["action"], purpose=row["purpose"], target=json.loads(row["target_json"]), not_before=row["not_before"], expires_at=row["expires_at"], required_count=row["required_count"], superseded_by=row["superseded_by"], consumed_digest=row["consumed_digest"], consumed_decision_id=row["consumed_decision_id"], consumed_at=row["consumed_at"], created_at=row["created_at"], updated_at=row["updated_at"], entries=entries, ) def _expire(self, approval_id: str, from_status: str) -> str: conn = self._conn() now = iso(self.now()) try: conn.execute("BEGIN IMMEDIATE") cur = conn.execute( "UPDATE approvals SET status='expired', updated_at=? " "WHERE id=? AND status=?", (now, approval_id, from_status), ) conn.commit() return "expired" if cur.rowcount == 1 else from_status except sqlite3.Error: conn.rollback() return from_status def add_entry( self, approval_id: str, subject_id: str, *, assurance: str | None = None, evidence_ref: str | None = None, ) -> Approval: if not subject_id: raise Unprocessable("subject_id is required") conn = self._conn() now = iso(self.now()) try: conn.execute("BEGIN IMMEDIATE") row = conn.execute("SELECT * FROM approvals WHERE id=?", (approval_id,)).fetchone() if row is None: conn.rollback() raise NotFound(approval_id) if row["status"] not in MUTABLE: conn.rollback() raise Conflict(f"cannot add entries in status {row['status']}") try: conn.execute( "INSERT INTO entries (approval_id, subject_id, approved_at, assurance, evidence_ref) " "VALUES (?,?,?,?,?)", (approval_id, subject_id, now, assurance, evidence_ref), ) except sqlite3.IntegrityError as exc: conn.rollback() raise DuplicateApprover(subject_id) from exc count = conn.execute( "SELECT COUNT(*) FROM entries WHERE approval_id=?", (approval_id,) ).fetchone()[0] became_approved = False if row["status"] == "requested" and count >= row["required_count"]: conn.execute( "UPDATE approvals SET status='approved', updated_at=? WHERE id=? AND status='requested'", (now, approval_id), ) became_approved = True else: conn.execute( "UPDATE approvals SET updated_at=? WHERE id=?", (now, approval_id) ) if became_approved: self._outbox_insert( conn, "issuance", approval_id, actor=row["actor"], extra={"binding_digest": row["binding_digest"], "required_count": row["required_count"]}, ) conn.commit() except sqlite3.Error as exc: conn.rollback() raise StoreUnavailable(str(exc)) from exc except Exception: conn.rollback() raise return self.get(approval_id) def revoke(self, approval_id: str) -> Approval: return self._terminal(approval_id, "revoked", "revocation") def supersede(self, approval_id: str, successor_id: str | None = None) -> dict[str, Any]: successor_id = successor_id or str(uuid4()) conn = self._conn() now = iso(self.now()) try: conn.execute("BEGIN IMMEDIATE") row = conn.execute("SELECT * FROM approvals WHERE id=?", (approval_id,)).fetchone() if row is None: conn.rollback() raise NotFound(approval_id) if row["status"] not in MUTABLE: conn.rollback() raise Conflict(f"cannot supersede in status {row['status']}") cur = conn.execute( "UPDATE approvals SET status='superseded', superseded_by=?, updated_at=? " "WHERE id=? AND status IN ('requested','approved')", (successor_id, now, approval_id), ) if cur.rowcount != 1: conn.rollback() raise Conflict("supersession lost the compare-and-swap") existing = conn.execute("SELECT id FROM approvals WHERE id=?", (successor_id,)).fetchone() created_successor = False if existing is None: conn.execute( """INSERT INTO approvals ( id, status, binding_json, binding_digest, pdp_digest, actor, principal, action, purpose, target_json, not_before, expires_at, required_count, superseded_by, created_at, updated_at ) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)""", ( successor_id, "requested", row["binding_json"], row["binding_digest"], row["pdp_digest"], row["actor"], row["principal"], row["action"], row["purpose"], row["target_json"], row["not_before"], row["expires_at"], row["required_count"], None, now, now, ), ) created_successor = True self._outbox_insert( conn, "supersession", approval_id, actor=row["actor"], extra={ "binding_digest": row["binding_digest"], "superseded_by": successor_id, }, ) conn.commit() except sqlite3.Error as exc: conn.rollback() raise StoreUnavailable(str(exc)) from exc except Exception: conn.rollback() raise return { "superseded": self.get(approval_id).as_dict(), "successor_id": successor_id, "successor_created": created_successor, } def _terminal(self, approval_id: str, status: str, event_class: str) -> Approval: conn = self._conn() now = iso(self.now()) try: conn.execute("BEGIN IMMEDIATE") row = conn.execute("SELECT * FROM approvals WHERE id=?", (approval_id,)).fetchone() if row is None: conn.rollback() raise NotFound(approval_id) cur = conn.execute( "UPDATE approvals SET status=?, updated_at=? " "WHERE id=? AND status IN ('requested','approved')", (status, now, approval_id), ) if cur.rowcount != 1: conn.rollback() raise Conflict(f"cannot {status} from status {row['status']}") self._outbox_insert( conn, event_class, approval_id, actor=row["actor"], extra={"binding_digest": row["binding_digest"]}, ) conn.commit() except sqlite3.Error as exc: conn.rollback() raise StoreUnavailable(str(exc)) from exc except Exception: conn.rollback() raise return self.get(approval_id) def consume( self, approval_id: str, request_digest: str | None, *, decision_id: str | None = None, ) -> dict[str, Any]: """Atomically spend an approval for one decision-bound request. Repeating the same request digest is an idempotent success. A different digest against a consumed object conflicts and the caller must not act. """ if request_digest is None: raise Unprocessable("request_digest is required") if not isinstance(request_digest, str): raise Unprocessable("request_digest must be a string") digest = require_digest(request_digest) if decision_id is not None and ( not isinstance(decision_id, str) or not decision_id ): raise Unprocessable("decision_id must be a non-empty string") conn = self._conn() now = iso(self.now()) idempotent = False try: conn.execute("BEGIN IMMEDIATE") row = conn.execute("SELECT * FROM approvals WHERE id=?", (approval_id,)).fetchone() if row is None: conn.rollback() raise NotFound(approval_id) if row["status"] == "consumed": if row["consumed_digest"] != digest: conn.rollback() raise Conflict("approval already consumed for a different request digest") idempotent = True stored_decision_id = row["consumed_decision_id"] conn.commit() return { "approval_id": approval_id, "status": "consumed", "request_digest": digest, **({"decision_id": stored_decision_id} if stored_decision_id else {}), "consumed_at": row["consumed_at"], "idempotent": idempotent, } if row["status"] != "approved": conn.rollback() raise Conflict(f"cannot consume from status {row['status']}") if self.now() < parse_iso(row["not_before"]) or self.now() >= parse_iso(row["expires_at"]): conn.rollback() raise Conflict("cannot consume outside validity window") cur = conn.execute( "UPDATE approvals SET status='consumed', consumed_digest=?, " "consumed_decision_id=?, consumed_at=?, updated_at=? " "WHERE id=? AND status='approved'", (digest, decision_id, now, now, approval_id), ) if cur.rowcount != 1: conn.rollback() raise Conflict("consumption lost the compare-and-swap") self._outbox_insert( conn, "use", approval_id, actor=row["actor"], extra={ "binding_digest": row["binding_digest"], "request_digest": digest, **({"decision_id": decision_id} if decision_id else {}), }, ) conn.commit() except sqlite3.Error as exc: conn.rollback() raise StoreUnavailable(str(exc)) from exc except Exception: conn.rollback() raise return { "approval_id": approval_id, "status": "consumed", "request_digest": digest, **({"decision_id": decision_id} if decision_id else {}), "consumed_at": now, "idempotent": idempotent, } def _outbox_insert( self, conn: sqlite3.Connection, event_class: str, approval_id: str | None, *, actor: str | None, extra: dict[str, Any] | None = None, ) -> str: if self.fail_outbox: raise StoreUnavailable("outbox insert failed") if event_class not in EVENT_CLASSES: raise Unprocessable(f"unknown event class {event_class}") event_id = str(uuid4()) created = iso(self.now()) resource = f"approval:{approval_id}" if approval_id else "approval-engine:heartbeat" details = {"class": event_class, **(extra or {})} if approval_id: details["approval_id"] = approval_id payload = { "schema_version": AUDIT_SCHEMA, "event_id": event_id, "observed_at": created, "tenant": self.tenant, "scope": SCOPE, "source": SOURCE, "actor": actor, "action": f"approval.{event_class}", "resource": resource, "outcome": "success", "reason": None, "details": details, } conn.execute( "INSERT INTO outbox (event_id, class, approval_id, payload_json, created_at, drained_at) " "VALUES (?,?,?,?,?,NULL)", (event_id, event_class, approval_id, json.dumps(payload, sort_keys=True), created), ) return event_id def claim(self, approval_id: str) -> dict[str, Any]: obj = self.get(approval_id) observed = self.now() not_after = observed + timedelta(seconds=self.freshness_ttl) state, valid_now, consumed, reason = self._evaluate(obj, observed) binding = { **canonical_binding(obj.binding), "digest": obj.binding_digest, } if obj.pdp_digest: binding["pdp_digest"] = obj.pdp_digest return { "schema_version": CLAIM_SCHEMA, "kind": "approval-claim", "yields_to": YIELDS_TO, "issuer": ISSUER, "approval_id": obj.id, "state": state, "valid_now": valid_now, "consumed": consumed, "binding": binding, "freshness": { "observed_at": iso(observed), "ttl_seconds": self.freshness_ttl, "not_after": iso(not_after), }, "validity": {"not_before": obj.not_before, "expires_at": obj.expires_at}, "reason_code": reason, } def _evaluate(self, obj: Approval, now: datetime) -> tuple[str, bool, bool, str]: consumed = obj.status == "consumed" if obj.status == "revoked": return "revoked", False, consumed, "revoked" if obj.status == "superseded": return "superseded", False, consumed, "superseded" if obj.status == "expired" or now >= parse_iso(obj.expires_at): return "expired", False, consumed, "expired" if obj.status == "consumed": return "consumed", False, True, "consumed" if obj.status == "requested": return "requested", False, False, "requested" if obj.status == "approved": if now < parse_iso(obj.not_before): return "approved", False, False, "not_yet_valid" if len({e["subject_id"] for e in obj.entries}) < obj.required_count: return "approved", False, False, "insufficient_approvers" return "valid", True, False, "ok" return obj.status, False, consumed, obj.status def emit_heartbeat(self) -> dict[str, Any]: conn = self._conn() counts = self.transition_counts() counts["heartbeat"] += 1 try: conn.execute("BEGIN IMMEDIATE") event_id = self._outbox_insert( conn, "heartbeat", None, actor=None, extra={"assertion": "nothing-to-report", "counts": counts}, ) conn.commit() except sqlite3.Error as exc: conn.rollback() raise StoreUnavailable(str(exc)) from exc except Exception: conn.rollback() raise return {"event_id": event_id, "assertion": "nothing-to-report", "counts": counts} def transition_counts(self) -> dict[str, int]: conn = self._conn() try: rows = conn.execute("SELECT class, COUNT(*) AS n FROM outbox GROUP BY class").fetchall() except sqlite3.Error as exc: raise StoreUnavailable(str(exc)) from exc counts = {c: 0 for c in EVENT_CLASSES} for row in rows: counts[row["class"]] = row["n"] return counts def outbox_stats(self) -> dict[str, Any]: conn = self._conn() now = self.now() try: rows = conn.execute( "SELECT event_id, class, created_at, drained_at FROM outbox" ).fetchall() except sqlite3.Error as exc: raise StoreUnavailable(str(exc)) from exc pending = [r for r in rows if r["drained_at"] is None] ages = [(now - parse_iso(r["created_at"])).total_seconds() for r in pending] return { "total": len(rows), "pending": len(pending), "max_age_seconds": max(ages) if ages else 0, "counts": self.transition_counts(), } def undrained(self) -> list[dict[str, Any]]: conn = self._conn() try: rows = conn.execute( "SELECT event_id, class, approval_id, payload_json, created_at " "FROM outbox WHERE drained_at IS NULL ORDER BY created_at" ).fetchall() except sqlite3.Error as exc: raise StoreUnavailable(str(exc)) from exc return [ { "event_id": r["event_id"], "class": r["class"], "approval_id": r["approval_id"], "payload": json.loads(r["payload_json"]), "created_at": r["created_at"], } for r in rows ] def drain(self, sink: Callable[[dict[str, Any]], None]) -> dict[str, int]: """Deliver undrained payloads. Sink failure leaves the row pending. Object mutations are already committed; this must not roll them back. """ delivered = 0 failed = 0 conn = self._conn() for item in self.undrained(): try: sink(item["payload"]) except Exception: failed += 1 continue try: conn.execute( "UPDATE outbox SET drained_at=? WHERE event_id=? AND drained_at IS NULL", (iso(self.now()), item["event_id"]), ) conn.commit() delivered += 1 except sqlite3.Error as exc: conn.rollback() raise StoreUnavailable(str(exc)) from exc return {"delivered": delivered, "failed": failed}