| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333 |
- """Governance audit normalization and tamper-evident evidence sealing."""
- from __future__ import annotations
- import hashlib
- import hmac
- import json
- from datetime import UTC, datetime
- from typing import Any
- from app.core.common.identifiers import new_governance_uid
- from app.core.data_source.redaction import redact_mapping
- AUDIT_CATEGORIES = (
- "authentication",
- "ingestion",
- "entity_resolution",
- "publication",
- "remediation",
- "knowledge_query",
- "authorization",
- "workflow_task",
- "data_product",
- "agent",
- "security_governance",
- )
- MAX_SEAL_EVENTS = 50_000
- class GovernanceAuditInvalid(ValueError):
- """Raised when an audit query or seal request is invalid."""
- class GovernanceAuditNotFound(LookupError):
- """Raised when a requested evidence seal does not exist."""
- def _utc(value: datetime, label: str) -> datetime:
- if not isinstance(value, datetime):
- raise GovernanceAuditInvalid(f"{label} must be a datetime")
- if value.tzinfo is None:
- raise GovernanceAuditInvalid(f"{label} must include a timezone")
- return value.astimezone(UTC)
- def _iso(value: datetime, label: str) -> str:
- return _utc(value, label).isoformat().replace("+00:00", "Z")
- def _canonical(value: Any) -> bytes:
- return json.dumps(
- value,
- ensure_ascii=False,
- sort_keys=True,
- separators=(",", ":"),
- ).encode("utf-8")
- def _digest(value: Any) -> str:
- return hashlib.sha256(_canonical(value)).hexdigest()
- def _required_text(value: Any, label: str, maximum: int) -> str:
- normalized = str(value or "").strip()
- if not normalized:
- raise GovernanceAuditInvalid(f"{label} is required")
- if len(normalized) > maximum:
- raise GovernanceAuditInvalid(f"{label} exceeds {maximum} characters")
- return normalized
- def _normalize_categories(categories=None) -> tuple[str, ...]:
- if categories is None:
- return AUDIT_CATEGORIES
- if isinstance(categories, str):
- categories = [categories]
- normalized = tuple(
- dict.fromkeys(str(item).strip() for item in categories if str(item).strip())
- )
- if not normalized:
- raise GovernanceAuditInvalid("at least one category is required")
- unsupported = sorted(set(normalized) - set(AUDIT_CATEGORIES))
- if unsupported:
- raise GovernanceAuditInvalid(
- f"unsupported category: {', '.join(unsupported)}"
- )
- return tuple(item for item in AUDIT_CATEGORIES if item in normalized)
- def _normalize_window(period_start, period_end) -> tuple[datetime, datetime]:
- start = _utc(period_start, "period_start")
- end = _utc(period_end, "period_end")
- if end <= start:
- raise GovernanceAuditInvalid("period_end must be after period_start")
- return start, end
- def _normalize_event(event: dict[str, Any]) -> dict[str, Any]:
- category = _required_text(event.get("category"), "category", 40)
- if category not in AUDIT_CATEGORIES:
- raise GovernanceAuditInvalid(f"unsupported category: {category}")
- occurred_at = event.get("occurred_at")
- return {
- "event_uid": _required_text(event.get("event_uid"), "event_uid", 160),
- "category": category,
- "action": _required_text(event.get("action"), "action", 80),
- "status": _required_text(event.get("status"), "status", 40),
- "actor_uid": str(event.get("actor_uid") or "")[:160] or None,
- "resource_type": _required_text(
- event.get("resource_type"), "resource_type", 80
- ),
- "resource_uid": str(event.get("resource_uid") or "")[:200] or None,
- "occurred_at": _iso(occurred_at, "occurred_at"),
- "safe_detail": redact_mapping(dict(event.get("safe_detail") or {})),
- }
- def _event_sort_key(event: dict[str, Any]) -> tuple[str, str]:
- return event["occurred_at"], event["event_uid"]
- def _root(events: list[dict[str, Any]]) -> str:
- ordered = sorted(events, key=_event_sort_key)
- return _digest([_digest(event) for event in ordered])
- def _signature_payload(seal: dict[str, Any]) -> dict[str, Any]:
- return {
- "uid": seal["uid"],
- "period_start": seal["period_start"],
- "period_end": seal["period_end"],
- "categories": list(seal["categories"]),
- "event_count": int(seal["event_count"]),
- "root_hash": seal["root_hash"],
- "key_version": seal["key_version"],
- "sealed_by": seal["sealed_by"],
- }
- class GovernanceAuditService:
- """Normalize distributed audit evidence and create signed evidence seals."""
- def __init__(
- self,
- repository,
- *,
- evidence_secret: str,
- key_version: str,
- uid_factory=new_governance_uid,
- now_factory=None,
- ):
- secret = str(evidence_secret or "").encode("utf-8")
- if len(secret) < 32:
- raise GovernanceAuditInvalid(
- "audit evidence secret must contain at least 32 bytes"
- )
- self.repository = repository
- self._secret = secret
- self.key_version = _required_text(key_version, "key_version", 80)
- self.uid_factory = uid_factory
- self.now_factory = now_factory or (lambda: datetime.now(UTC))
- def _events(self, *, categories, period_start, period_end):
- normalized_categories = _normalize_categories(categories)
- start, end = _normalize_window(period_start, period_end)
- events = [
- _normalize_event(event)
- for event in self.repository.fetch_events(
- categories=normalized_categories,
- period_start=start,
- period_end=end,
- )
- ]
- return normalized_categories, start, end, events
- def list_events(
- self,
- *,
- period_start,
- period_end,
- categories=None,
- action=None,
- status=None,
- page=1,
- page_size=20,
- ) -> dict[str, Any]:
- try:
- page = int(page)
- page_size = int(page_size)
- except (TypeError, ValueError) as exc:
- raise GovernanceAuditInvalid(
- "page and page_size must be integers"
- ) from exc
- if page < 1:
- raise GovernanceAuditInvalid("page must be at least 1")
- if page_size < 1 or page_size > 100:
- raise GovernanceAuditInvalid("page_size must be between 1 and 100")
- selected, start, end, events = self._events(
- categories=categories,
- period_start=period_start,
- period_end=period_end,
- )
- if action:
- events = [event for event in events if event["action"] == action]
- if status:
- events = [event for event in events if event["status"] == status]
- events.sort(key=_event_sort_key, reverse=True)
- offset = (page - 1) * page_size
- return {
- "period_start": _iso(start, "period_start"),
- "period_end": _iso(end, "period_end"),
- "categories": list(selected),
- "records": events[offset : offset + page_size],
- "page": page,
- "page_size": page_size,
- "total": len(events),
- }
- def coverage(self, *, period_start, period_end) -> dict[str, Any]:
- _, start, end, events = self._events(
- categories=AUDIT_CATEGORIES,
- period_start=period_start,
- period_end=period_end,
- )
- counts = dict.fromkeys(AUDIT_CATEGORIES, 0)
- latest = dict.fromkeys(AUDIT_CATEGORIES)
- for event in sorted(events, key=_event_sort_key):
- category = event["category"]
- counts[category] += 1
- latest[category] = event["occurred_at"]
- return {
- "period_start": _iso(start, "period_start"),
- "period_end": _iso(end, "period_end"),
- "categories": [
- {
- "category": category,
- "count": counts[category],
- "latest_at": latest[category],
- "available": counts[category] > 0,
- }
- for category in AUDIT_CATEGORIES
- ],
- }
- def _sign(self, seal: dict[str, Any]) -> str:
- return hmac.new(
- self._secret,
- _canonical(_signature_payload(seal)),
- hashlib.sha256,
- ).hexdigest()
- def create_seal(
- self,
- *,
- period_start,
- period_end,
- actor_uid,
- categories=None,
- ) -> dict[str, Any]:
- selected, start, end, events = self._events(
- categories=categories,
- period_start=period_start,
- period_end=period_end,
- )
- if end > _utc(self.now_factory(), "current_time"):
- raise GovernanceAuditInvalid(
- "evidence seal period_end cannot be in the future"
- )
- if len(events) > MAX_SEAL_EVENTS:
- raise GovernanceAuditInvalid(
- "evidence seal cannot contain more than 50,000 events"
- )
- seal = {
- "uid": self.uid_factory(),
- "period_start": _iso(start, "period_start"),
- "period_end": _iso(end, "period_end"),
- "categories": list(selected),
- "event_count": len(events),
- "root_hash": _root(events),
- "key_version": self.key_version,
- "sealed_by": _required_text(actor_uid, "actor_uid", 160),
- }
- seal["signature"] = self._sign(seal)
- return self.repository.save_seal(seal)
- def verify_seal(self, seal_uid) -> dict[str, Any]:
- seal = self.repository.get_seal(
- _required_text(seal_uid, "seal_uid", 160)
- )
- if seal is None:
- raise GovernanceAuditNotFound("evidence seal not found")
- expected_signature = self._sign(seal)
- if not hmac.compare_digest(
- expected_signature, str(seal.get("signature") or "")
- ):
- return {
- **seal,
- "integrity_status": "invalid_signature",
- "actual_event_count": None,
- "actual_root_hash": None,
- }
- categories, _, _, events = self._events(
- categories=seal["categories"],
- period_start=datetime.fromisoformat(
- seal["period_start"].replace("Z", "+00:00")
- ),
- period_end=datetime.fromisoformat(
- seal["period_end"].replace("Z", "+00:00")
- ),
- )
- actual_root = _root(events)
- intact = (
- list(categories) == list(seal["categories"])
- and len(events) == int(seal["event_count"])
- and hmac.compare_digest(actual_root, seal["root_hash"])
- )
- return {
- **seal,
- "integrity_status": "intact" if intact else "tampered",
- "actual_event_count": len(events),
- "actual_root_hash": actual_root,
- }
- def list_seals(self, *, limit=50):
- try:
- limit = int(limit)
- except (TypeError, ValueError) as exc:
- raise GovernanceAuditInvalid("limit must be an integer") from exc
- if limit < 1 or limit > 100:
- raise GovernanceAuditInvalid("limit must be between 1 and 100")
- return list(self.repository.list_seals(limit=limit))
|