governance_audit.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328
  1. """Governance audit normalization and tamper-evident evidence sealing."""
  2. from __future__ import annotations
  3. import hashlib
  4. import hmac
  5. import json
  6. from datetime import UTC, datetime
  7. from typing import Any
  8. from app.core.common.identifiers import new_governance_uid
  9. from app.core.data_source.redaction import redact_mapping
  10. AUDIT_CATEGORIES = (
  11. "authentication",
  12. "ingestion",
  13. "entity_resolution",
  14. "publication",
  15. "remediation",
  16. "knowledge_query",
  17. )
  18. MAX_SEAL_EVENTS = 50_000
  19. class GovernanceAuditInvalid(ValueError):
  20. """Raised when an audit query or seal request is invalid."""
  21. class GovernanceAuditNotFound(LookupError):
  22. """Raised when a requested evidence seal does not exist."""
  23. def _utc(value: datetime, label: str) -> datetime:
  24. if not isinstance(value, datetime):
  25. raise GovernanceAuditInvalid(f"{label} must be a datetime")
  26. if value.tzinfo is None:
  27. raise GovernanceAuditInvalid(f"{label} must include a timezone")
  28. return value.astimezone(UTC)
  29. def _iso(value: datetime, label: str) -> str:
  30. return _utc(value, label).isoformat().replace("+00:00", "Z")
  31. def _canonical(value: Any) -> bytes:
  32. return json.dumps(
  33. value,
  34. ensure_ascii=False,
  35. sort_keys=True,
  36. separators=(",", ":"),
  37. ).encode("utf-8")
  38. def _digest(value: Any) -> str:
  39. return hashlib.sha256(_canonical(value)).hexdigest()
  40. def _required_text(value: Any, label: str, maximum: int) -> str:
  41. normalized = str(value or "").strip()
  42. if not normalized:
  43. raise GovernanceAuditInvalid(f"{label} is required")
  44. if len(normalized) > maximum:
  45. raise GovernanceAuditInvalid(f"{label} exceeds {maximum} characters")
  46. return normalized
  47. def _normalize_categories(categories=None) -> tuple[str, ...]:
  48. if categories is None:
  49. return AUDIT_CATEGORIES
  50. if isinstance(categories, str):
  51. categories = [categories]
  52. normalized = tuple(
  53. dict.fromkeys(str(item).strip() for item in categories if str(item).strip())
  54. )
  55. if not normalized:
  56. raise GovernanceAuditInvalid("at least one category is required")
  57. unsupported = sorted(set(normalized) - set(AUDIT_CATEGORIES))
  58. if unsupported:
  59. raise GovernanceAuditInvalid(
  60. f"unsupported category: {', '.join(unsupported)}"
  61. )
  62. return tuple(item for item in AUDIT_CATEGORIES if item in normalized)
  63. def _normalize_window(period_start, period_end) -> tuple[datetime, datetime]:
  64. start = _utc(period_start, "period_start")
  65. end = _utc(period_end, "period_end")
  66. if end <= start:
  67. raise GovernanceAuditInvalid("period_end must be after period_start")
  68. return start, end
  69. def _normalize_event(event: dict[str, Any]) -> dict[str, Any]:
  70. category = _required_text(event.get("category"), "category", 40)
  71. if category not in AUDIT_CATEGORIES:
  72. raise GovernanceAuditInvalid(f"unsupported category: {category}")
  73. occurred_at = event.get("occurred_at")
  74. return {
  75. "event_uid": _required_text(event.get("event_uid"), "event_uid", 160),
  76. "category": category,
  77. "action": _required_text(event.get("action"), "action", 80),
  78. "status": _required_text(event.get("status"), "status", 40),
  79. "actor_uid": str(event.get("actor_uid") or "")[:160] or None,
  80. "resource_type": _required_text(
  81. event.get("resource_type"), "resource_type", 80
  82. ),
  83. "resource_uid": str(event.get("resource_uid") or "")[:200] or None,
  84. "occurred_at": _iso(occurred_at, "occurred_at"),
  85. "safe_detail": redact_mapping(dict(event.get("safe_detail") or {})),
  86. }
  87. def _event_sort_key(event: dict[str, Any]) -> tuple[str, str]:
  88. return event["occurred_at"], event["event_uid"]
  89. def _root(events: list[dict[str, Any]]) -> str:
  90. ordered = sorted(events, key=_event_sort_key)
  91. return _digest([_digest(event) for event in ordered])
  92. def _signature_payload(seal: dict[str, Any]) -> dict[str, Any]:
  93. return {
  94. "uid": seal["uid"],
  95. "period_start": seal["period_start"],
  96. "period_end": seal["period_end"],
  97. "categories": list(seal["categories"]),
  98. "event_count": int(seal["event_count"]),
  99. "root_hash": seal["root_hash"],
  100. "key_version": seal["key_version"],
  101. "sealed_by": seal["sealed_by"],
  102. }
  103. class GovernanceAuditService:
  104. """Normalize distributed audit evidence and create signed evidence seals."""
  105. def __init__(
  106. self,
  107. repository,
  108. *,
  109. evidence_secret: str,
  110. key_version: str,
  111. uid_factory=new_governance_uid,
  112. now_factory=None,
  113. ):
  114. secret = str(evidence_secret or "").encode("utf-8")
  115. if len(secret) < 32:
  116. raise GovernanceAuditInvalid(
  117. "audit evidence secret must contain at least 32 bytes"
  118. )
  119. self.repository = repository
  120. self._secret = secret
  121. self.key_version = _required_text(key_version, "key_version", 80)
  122. self.uid_factory = uid_factory
  123. self.now_factory = now_factory or (lambda: datetime.now(UTC))
  124. def _events(self, *, categories, period_start, period_end):
  125. normalized_categories = _normalize_categories(categories)
  126. start, end = _normalize_window(period_start, period_end)
  127. events = [
  128. _normalize_event(event)
  129. for event in self.repository.fetch_events(
  130. categories=normalized_categories,
  131. period_start=start,
  132. period_end=end,
  133. )
  134. ]
  135. return normalized_categories, start, end, events
  136. def list_events(
  137. self,
  138. *,
  139. period_start,
  140. period_end,
  141. categories=None,
  142. action=None,
  143. status=None,
  144. page=1,
  145. page_size=20,
  146. ) -> dict[str, Any]:
  147. try:
  148. page = int(page)
  149. page_size = int(page_size)
  150. except (TypeError, ValueError) as exc:
  151. raise GovernanceAuditInvalid(
  152. "page and page_size must be integers"
  153. ) from exc
  154. if page < 1:
  155. raise GovernanceAuditInvalid("page must be at least 1")
  156. if page_size < 1 or page_size > 100:
  157. raise GovernanceAuditInvalid("page_size must be between 1 and 100")
  158. selected, start, end, events = self._events(
  159. categories=categories,
  160. period_start=period_start,
  161. period_end=period_end,
  162. )
  163. if action:
  164. events = [event for event in events if event["action"] == action]
  165. if status:
  166. events = [event for event in events if event["status"] == status]
  167. events.sort(key=_event_sort_key, reverse=True)
  168. offset = (page - 1) * page_size
  169. return {
  170. "period_start": _iso(start, "period_start"),
  171. "period_end": _iso(end, "period_end"),
  172. "categories": list(selected),
  173. "records": events[offset : offset + page_size],
  174. "page": page,
  175. "page_size": page_size,
  176. "total": len(events),
  177. }
  178. def coverage(self, *, period_start, period_end) -> dict[str, Any]:
  179. _, start, end, events = self._events(
  180. categories=AUDIT_CATEGORIES,
  181. period_start=period_start,
  182. period_end=period_end,
  183. )
  184. counts = dict.fromkeys(AUDIT_CATEGORIES, 0)
  185. latest = dict.fromkeys(AUDIT_CATEGORIES)
  186. for event in sorted(events, key=_event_sort_key):
  187. category = event["category"]
  188. counts[category] += 1
  189. latest[category] = event["occurred_at"]
  190. return {
  191. "period_start": _iso(start, "period_start"),
  192. "period_end": _iso(end, "period_end"),
  193. "categories": [
  194. {
  195. "category": category,
  196. "count": counts[category],
  197. "latest_at": latest[category],
  198. "available": counts[category] > 0,
  199. }
  200. for category in AUDIT_CATEGORIES
  201. ],
  202. }
  203. def _sign(self, seal: dict[str, Any]) -> str:
  204. return hmac.new(
  205. self._secret,
  206. _canonical(_signature_payload(seal)),
  207. hashlib.sha256,
  208. ).hexdigest()
  209. def create_seal(
  210. self,
  211. *,
  212. period_start,
  213. period_end,
  214. actor_uid,
  215. categories=None,
  216. ) -> dict[str, Any]:
  217. selected, start, end, events = self._events(
  218. categories=categories,
  219. period_start=period_start,
  220. period_end=period_end,
  221. )
  222. if end > _utc(self.now_factory(), "current_time"):
  223. raise GovernanceAuditInvalid(
  224. "evidence seal period_end cannot be in the future"
  225. )
  226. if len(events) > MAX_SEAL_EVENTS:
  227. raise GovernanceAuditInvalid(
  228. "evidence seal cannot contain more than 50,000 events"
  229. )
  230. seal = {
  231. "uid": self.uid_factory(),
  232. "period_start": _iso(start, "period_start"),
  233. "period_end": _iso(end, "period_end"),
  234. "categories": list(selected),
  235. "event_count": len(events),
  236. "root_hash": _root(events),
  237. "key_version": self.key_version,
  238. "sealed_by": _required_text(actor_uid, "actor_uid", 160),
  239. }
  240. seal["signature"] = self._sign(seal)
  241. return self.repository.save_seal(seal)
  242. def verify_seal(self, seal_uid) -> dict[str, Any]:
  243. seal = self.repository.get_seal(
  244. _required_text(seal_uid, "seal_uid", 160)
  245. )
  246. if seal is None:
  247. raise GovernanceAuditNotFound("evidence seal not found")
  248. expected_signature = self._sign(seal)
  249. if not hmac.compare_digest(
  250. expected_signature, str(seal.get("signature") or "")
  251. ):
  252. return {
  253. **seal,
  254. "integrity_status": "invalid_signature",
  255. "actual_event_count": None,
  256. "actual_root_hash": None,
  257. }
  258. categories, _, _, events = self._events(
  259. categories=seal["categories"],
  260. period_start=datetime.fromisoformat(
  261. seal["period_start"].replace("Z", "+00:00")
  262. ),
  263. period_end=datetime.fromisoformat(
  264. seal["period_end"].replace("Z", "+00:00")
  265. ),
  266. )
  267. actual_root = _root(events)
  268. intact = (
  269. list(categories) == list(seal["categories"])
  270. and len(events) == int(seal["event_count"])
  271. and hmac.compare_digest(actual_root, seal["root_hash"])
  272. )
  273. return {
  274. **seal,
  275. "integrity_status": "intact" if intact else "tampered",
  276. "actual_event_count": len(events),
  277. "actual_root_hash": actual_root,
  278. }
  279. def list_seals(self, *, limit=50):
  280. try:
  281. limit = int(limit)
  282. except (TypeError, ValueError) as exc:
  283. raise GovernanceAuditInvalid("limit must be an integer") from exc
  284. if limit < 1 or limit > 100:
  285. raise GovernanceAuditInvalid("limit must be between 1 and 100")
  286. return list(self.repository.list_seals(limit=limit))