"""Synchronous authorization custody using Audit Core's eight-field contract.""" from __future__ import annotations import asyncio import json from datetime import datetime, timezone from pathlib import Path from uuid import uuid4 import httpx from hub_core.security.identity import AccessFailure, require_https class AuditCoreSink: """No allow returns before durable remote custody acknowledges its record. A lost receipt blocks execution, even if the archive already stored the attempt. Committed outcomes arrive separately from the transactional outbox. """ def __init__(self, *, base_url: str, token_file: Path, client: httpx.AsyncClient): require_https(base_url) self.base_url = base_url.rstrip("/") self.token_file, self.client = token_file, client async def readiness(self) -> None: response = await self.client.get(self.base_url + "/readyz", timeout=2, follow_redirects=False) response.raise_for_status() receipt = response.json() if (receipt.get("status") != "ok" or receipt.get("durable") is not True or receipt.get("custody_class") not in {"archive", "operational"}): raise ValueError("operational audit custody is required") async def append(self, record: dict) -> None: correlation, outcome = record.get("correlation_id"), record.get("outcome") if not isinstance(correlation, str) or not correlation or outcome not in {"authorized", "denied", "refused"}: raise AccessFailure(503, "audit_unavailable") await self._deliver({ "id": str(uuid4()), "type": "hub.access." + outcome, "source": "hub-core", "subject": "hub-access:" + correlation, "tenant": "tenant:platform", "correlation_id": correlation, "occurred_at": datetime.now(timezone.utc).isoformat(), "data": record, }) async def append_outcome(self, event: dict) -> None: # The transaction supplies the immutable ID, timestamp and full envelope. if (event.get("type") != "hub.operation.committed" or event.get("source") != "hub-core" or event.get("tenant") != "tenant:platform" or not isinstance(event.get("id"), str) or not event["id"] or set(event) != {"id", "type", "source", "subject", "tenant", "correlation_id", "occurred_at", "data"}): raise AccessFailure(503, "audit_unavailable") await self._deliver(event) async def _deliver(self, event: dict) -> None: try: async with asyncio.timeout(3): await self.readiness() token = self.token_file.read_text().strip() if not token or not token.isascii() or any(c.isspace() for c in token): raise ValueError("invalid sender credential") raw = json.dumps(event, ensure_ascii=False, allow_nan=False).encode() if len(raw) > 256 * 1024: raise ValueError("audit envelope exceeds receiver limit") response = await self.client.post( self.base_url + "/v1/events", content=raw, timeout=2, follow_redirects=False, headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json", "Idempotency-Key": event["id"]}, ) expected = {200: "duplicate", 202: "accepted"}.get(response.status_code) receipt = response.json() if (expected is None or receipt.get("status") != expected or not isinstance(receipt.get("reference"), str) or not receipt["reference"]): raise ValueError("audit custody not acknowledged") except Exception as exc: # Never return receiver bodies, credential values or private paths. raise AccessFailure(503, "audit_unavailable") from exc