trusted_delivery.py 42 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670
  1. """Closed, default-deny contracts for fine-grained trusted delivery.
  2. This module intentionally never executes a data query or contacts a target.
  3. It produces a bounded authorization projection and digest-only provider envelope.
  4. Concrete enterprise providers remain external integration work.
  5. """
  6. from __future__ import annotations
  7. import copy
  8. import hashlib
  9. import json
  10. import re
  11. import uuid
  12. from collections.abc import Callable, Mapping
  13. from datetime import UTC, datetime
  14. from typing import Any
  15. from app.core.common.identifiers import new_governance_uid
  16. from app.core.common.timezone_utils import now_china
  17. _PROVIDERS = frozenset({"database", "api", "file", "iam", "kms", "dlp", "siem"})
  18. _ENVIRONMENTS = frozenset({"development", "test", "production"})
  19. _ACTIONS = frozenset({"read", "use", "display", "export"})
  20. _CLASSIFICATIONS = frozenset({"public", "internal", "sensitive", "highly_sensitive"})
  21. _FIELD = re.compile(r"^[A-Za-z_][A-Za-z0-9_.-]{0,199}$")
  22. _CODE = re.compile(r"^[A-Z][A-Z0-9_]{2,119}$")
  23. _DIGEST = re.compile(r"^[0-9a-f]{64}$")
  24. _UNSAFE_TEXT = re.compile(
  25. r"(?i)(?:https?://|ftp://|file://|www\.|\b(?:secret|token|password|private[ _-]?key)\b|"
  26. r"-----BEGIN|<script|javascript:|;)")
  27. _UNSAFE_SQL = re.compile(r"(?i)\b(?:select|insert|update|delete|drop|alter)\s+")
  28. def _closed(value: Any, allowed: set[str], label: str) -> dict[str, Any]:
  29. if not isinstance(value, dict):
  30. raise ValueError(f"{label} must be an object")
  31. extra = sorted(set(value) - allowed)
  32. if extra:
  33. raise ValueError(f"{label} contains unsupported fields: {', '.join(extra)}")
  34. result = copy.deepcopy(value)
  35. _assert_safe_value(result, label)
  36. return result
  37. def _assert_safe_value(value: Any, label: str) -> None:
  38. """Reject values that could turn an authorization reference into payload data."""
  39. if isinstance(value, str) and _UNSAFE_TEXT.search(value):
  40. raise ValueError(f"{label} contains unsafe content")
  41. if isinstance(value, dict):
  42. for item in value.values():
  43. _assert_safe_value(item, label)
  44. elif isinstance(value, list):
  45. for item in value:
  46. _assert_safe_value(item, label)
  47. def _text(value: Any, label: str, maximum: int = 300) -> str:
  48. if not isinstance(value, str) or not value.strip():
  49. raise ValueError(f"{label} is required")
  50. result = value.strip()
  51. if len(result) > maximum:
  52. raise ValueError(f"{label} exceeds {maximum} characters")
  53. return result
  54. def _uid(value: Any, label: str) -> str:
  55. try:
  56. return str(uuid.UUID(str(value)))
  57. except (TypeError, ValueError, AttributeError) as exc:
  58. raise ValueError(f"{label} must be a UUID") from exc
  59. def _time(value: Any, label: str) -> datetime:
  60. try:
  61. result = value if isinstance(value, datetime) else datetime.fromisoformat(str(value).replace("Z", "+00:00"))
  62. except (TypeError, ValueError) as exc:
  63. raise ValueError(f"{label} must be ISO-8601") from exc
  64. if result.tzinfo is None:
  65. raise ValueError(f"{label} must include a timezone")
  66. return result.astimezone(UTC)
  67. def _values(value: Any, label: str, *, minimum: int = 1, maximum: int = 200) -> list[Any]:
  68. if not isinstance(value, list) or not minimum <= len(value) <= maximum:
  69. raise ValueError(f"{label} must contain between {minimum} and {maximum} items")
  70. return copy.deepcopy(value)
  71. def _digest(value: Any) -> str:
  72. return hashlib.sha256(json.dumps(value, sort_keys=True, separators=(",", ":")).encode()).hexdigest()
  73. def _digest_value(value: Any, label: str) -> str:
  74. result = _text(value, label, 64).lower()
  75. if not _DIGEST.fullmatch(result):
  76. raise ValueError(f"{label} must be SHA-256")
  77. return result
  78. def _field(value: Any, label: str) -> str:
  79. result = _text(value, label, 200)
  80. if not _FIELD.fullmatch(result):
  81. raise ValueError(f"{label} is invalid")
  82. return result
  83. def _simple_values(value: Any, label: str, allowed: frozenset[str]) -> list[str]:
  84. result = sorted({_text(item, label, 100) for item in _values(value, label)})
  85. if not set(result) <= allowed:
  86. raise ValueError(f"unsupported {label}")
  87. return result
  88. def _row_filter(value: Any) -> dict[str, Any] | None:
  89. if value is None:
  90. return None
  91. body = _closed(value, {"op", "field", "value", "items", "item"}, "row filter")
  92. op = _text(body.get("op"), "row filter op", 10)
  93. if op not in {"eq", "in", "and", "or", "not"}:
  94. raise ValueError("unsupported row filter op")
  95. if _UNSAFE_SQL.search(json.dumps(body, sort_keys=True)):
  96. raise ValueError("row filter contains unsafe content")
  97. if op in {"eq", "in"}:
  98. if set(body) - {"op", "field", "value"}:
  99. raise ValueError("row filter shape is invalid")
  100. field = _field(body.get("field"), "row filter field")
  101. raw = body.get("value")
  102. if op == "eq" and not isinstance(raw, (str, int, float, bool)):
  103. raise ValueError("row filter eq value must be scalar")
  104. if op == "in":
  105. if not isinstance(raw, list) or not 1 <= len(raw) <= 100 or not all(isinstance(item, (str, int, float, bool)) for item in raw):
  106. raise ValueError("row filter in value must be bounded scalar list")
  107. raw = sorted(raw, key=lambda item: (type(item).__name__, str(item)))
  108. return {"op": op, "field": field, "value": raw}
  109. if op == "not":
  110. if set(body) != {"op", "item"}:
  111. raise ValueError("row filter not shape is invalid")
  112. return {"op": op, "item": _row_filter(body["item"])}
  113. if set(body) != {"op", "items"}:
  114. raise ValueError("row filter boolean shape is invalid")
  115. items = [_row_filter(item) for item in _values(body.get("items"), "row filter items", minimum=1, maximum=20)]
  116. return {"op": op, "items": items}
  117. class ClosedProviderRegistry:
  118. """Provider names are a closed enum; tests may inject non-network adapters."""
  119. def __init__(self, adapters: Mapping[str, Any] | None = None, *, test_mode: bool = False):
  120. adapters = dict(adapters or {})
  121. if set(adapters) - _PROVIDERS:
  122. raise ValueError("unsupported trusted-delivery provider")
  123. if adapters and not test_mode:
  124. raise PermissionError("enterprise providers require approved external integration")
  125. self._adapters = adapters
  126. self._test_mode = test_mode
  127. @classmethod
  128. def for_tests(cls, adapters: Mapping[str, Any]):
  129. return cls(adapters, test_mode=True)
  130. @property
  131. def is_enterprise_ready(self) -> bool:
  132. return False
  133. def dependency_status(self) -> dict[str, str]:
  134. return dict.fromkeys(sorted(_PROVIDERS), "TBD_EXTERNAL")
  135. def get(self, provider: str):
  136. provider = _text(provider, "provider", 30)
  137. if provider not in _PROVIDERS:
  138. raise ValueError("unsupported trusted-delivery provider")
  139. adapter = self._adapters.get(provider)
  140. if adapter is None:
  141. raise PermissionError(f"trusted-delivery provider {provider} is disabled")
  142. return adapter
  143. class TrustedDeliveryService:
  144. """Durable-service façade. Repositories own the database transaction boundary."""
  145. def __init__(
  146. self, repository, *, provider_registry: ClosedProviderRegistry | None = None,
  147. uid_factory: Callable[[], str] = new_governance_uid,
  148. now_factory: Callable[[], datetime] = now_china,
  149. ):
  150. self.repository = repository
  151. self.provider_registry = provider_registry or ClosedProviderRegistry()
  152. self.uid_factory = uid_factory
  153. self.now_factory = now_factory
  154. def _actor(self, value: Any) -> str:
  155. actor = _uid(value, "actor_uid")
  156. if self.repository.users_available({actor}) != {actor}:
  157. raise ValueError("trusted-delivery actor is unavailable")
  158. return actor
  159. def _event(self, action: str, actor_uid: str, detail: dict[str, Any]) -> None:
  160. if hasattr(self.repository, "add_evidence"):
  161. self.repository.add_evidence({
  162. "uid": self.uid_factory(), "action": action, "actor_uid": actor_uid,
  163. "safe_detail": copy.deepcopy(detail), "created_at": self.now_factory().isoformat(),
  164. })
  165. def create_policy_version(self, payload: Any, *, actor_uid: str) -> dict[str, Any]:
  166. body = _closed(payload, {"code", "version", "selector", "resource_rules"}, "trusted delivery policy")
  167. actor = self._actor(actor_uid)
  168. code = _text(body.get("code"), "policy code", 120).upper()
  169. if not _CODE.fullmatch(code):
  170. raise ValueError("policy code is invalid")
  171. selector = _closed(body.get("selector"), {"subjects", "roles", "business_domains", "assets", "purposes", "environments", "actions", "classifications", "expires_at"}, "policy selector")
  172. subjects = sorted({_uid(item, "policy subject") for item in _values(selector.get("subjects", []), "policy subjects", minimum=0)})
  173. roles = sorted({_text(item, "policy role", 80) for item in _values(selector.get("roles", []), "policy roles", minimum=0)})
  174. if not subjects and not roles:
  175. raise ValueError("policy selector requires a subject or role")
  176. if subjects and self.repository.users_available(set(subjects)) != set(subjects):
  177. raise ValueError("policy selector contains unavailable users")
  178. rules = _closed(body.get("resource_rules"), {"row_filter", "field_rules"}, "resource rules")
  179. field_rules = []
  180. for item in _values(rules.get("field_rules"), "field rules", maximum=500):
  181. field_rule = _closed(item, {"field", "effect", "mask_ref"}, "field rule")
  182. effect = _text(field_rule.get("effect"), "field effect", 10)
  183. if effect not in {"allow", "deny"}:
  184. raise ValueError("unsupported field effect")
  185. mask_ref = field_rule.get("mask_ref")
  186. if mask_ref is not None:
  187. mask_ref = _text(mask_ref, "mask ref", 120)
  188. if not re.fullmatch(r"[a-z][a-z0-9-]{2,119}-v[1-9][0-9]*", mask_ref):
  189. raise ValueError("mask ref is invalid")
  190. field_rules.append({"field": _field(field_rule.get("field"), "field rule field"), "effect": effect, "mask_ref": mask_ref})
  191. if len({item["field"] for item in field_rules}) != len(field_rules):
  192. raise ValueError("field rules must not duplicate fields")
  193. expires_at = _time(selector.get("expires_at"), "policy expiry")
  194. now = self.now_factory().astimezone(UTC)
  195. if expires_at <= now:
  196. raise ValueError("policy expiry must be in the future")
  197. record = {
  198. "uid": self.uid_factory(), "code": code, "version": _text(body.get("version"), "policy version", 30),
  199. # A policy never becomes effective by creation alone. This keeps
  200. # both first and subsequent versions behind the durable approval
  201. # transition, rather than treating an empty policy history as an
  202. # implicit approval.
  203. "status": "draft", "selector": {
  204. "subjects": subjects, "roles": roles,
  205. "business_domains": sorted({_uid(item, "business domain") for item in _values(selector.get("business_domains"), "business domains")}),
  206. "assets": sorted({_uid(item, "asset") for item in _values(selector.get("assets"), "assets")}),
  207. "purposes": sorted({_text(item, "purpose", 100) for item in _values(selector.get("purposes"), "purposes")}),
  208. "environments": _simple_values(selector.get("environments"), "environments", _ENVIRONMENTS),
  209. "actions": _simple_values(selector.get("actions"), "actions", _ACTIONS),
  210. "classifications": _simple_values(selector.get("classifications"), "classifications", _CLASSIFICATIONS),
  211. "expires_at": expires_at.isoformat(),
  212. },
  213. "resource_rules": {"row_filter": _row_filter(rules.get("row_filter")), "field_rules": sorted(field_rules, key=lambda item: item["field"])},
  214. "created_by": actor, "created_at": now.isoformat(), "current_version": 1,
  215. }
  216. saved = self.repository.create_policy_version(record)
  217. self._event("policy_version_created", actor, {"policy_code": code, "version": record["version"]})
  218. return saved
  219. def _lifecycle_request(self, payload: Any, *, label: str) -> tuple[dict[str, Any], str, str, str]:
  220. body = _closed(payload, {"code", "version", "approval_ref", "approval_digest", "idempotency_key"}, label)
  221. actor = self._actor(payload.get("actor_uid")) if isinstance(payload, dict) and "actor_uid" in payload else None
  222. # Callers supply the actor separately. This branch only keeps direct
  223. # service use closed if a malformed payload is passed.
  224. if actor is not None:
  225. raise ValueError("actor_uid must not be in lifecycle payload")
  226. return (
  227. body,
  228. _text(body.get("code"), "policy code", 120).upper(),
  229. _text(body.get("version"), "policy version", 30),
  230. _text(body.get("idempotency_key"), "idempotency key", 160),
  231. )
  232. def _policy_transition(self, payload: Any, *, actor_uid: str, operation: str) -> dict[str, Any]:
  233. body = _closed(payload, {"code", "version", "approval_ref", "approval_digest", "idempotency_key"}, "policy transition")
  234. actor = self._actor(actor_uid)
  235. code = _text(body.get("code"), "policy code", 120).upper()
  236. if not _CODE.fullmatch(code):
  237. raise ValueError("policy code is invalid")
  238. version = _text(body.get("version"), "policy version", 30)
  239. approval_ref = _text(body.get("approval_ref"), "approval reference", 200)
  240. approval_digest = _digest_value(body.get("approval_digest"), "approval digest")
  241. key = _text(body.get("idempotency_key"), "idempotency key", 160)
  242. request_digest = _digest({"operation": operation, "code": code, "version": version, "approval_ref": approval_ref, "approval_digest": approval_digest})
  243. if hasattr(self.repository, "transition_policy"):
  244. saved = self.repository.transition_policy(code, version, operation, key, request_digest, actor, approval_ref, approval_digest)
  245. elif hasattr(self.repository, "policies"):
  246. # Small test repositories intentionally do not emulate PostgreSQL;
  247. # preserve the same state rules for service-level tests.
  248. history = getattr(self.repository, "policy_transitions", {})
  249. existing = history.get(key)
  250. if existing:
  251. if existing["request_digest"] != request_digest:
  252. raise RuntimeError("trusted delivery idempotency conflict")
  253. return copy.deepcopy(existing["result"])
  254. target = next((item for item in self.repository.policies.values() if item["code"] == code and item["version"] == version), None)
  255. if target is None:
  256. raise LookupError("trusted delivery policy version was not found")
  257. for item in self.repository.policies.values():
  258. if item["code"] == code and item["uid"] != target["uid"] and item["status"] == "active":
  259. item["status"] = "retired"
  260. item["current_version"] += 1
  261. if target["status"] != "active":
  262. target["status"] = "active"
  263. target["current_version"] += 1
  264. saved = {key: target[key] for key in ("uid", "code", "version", "status", "current_version")}
  265. history[key] = {"request_digest": request_digest, "result": copy.deepcopy(saved)}
  266. self.repository.policy_transitions = history
  267. else:
  268. raise RuntimeError("durable policy transition repository is required")
  269. self._event(f"policy_{operation}", actor, {"policy_code": code, "version": version, "approval_digest": approval_digest})
  270. return saved
  271. def activate_policy_version(self, payload: Any, *, actor_uid: str) -> dict[str, Any]:
  272. return self._policy_transition(payload, actor_uid=actor_uid, operation="activate")
  273. def rollback_policy_version(self, payload: Any, *, actor_uid: str) -> dict[str, Any]:
  274. return self._policy_transition(payload, actor_uid=actor_uid, operation="rollback")
  275. def evaluate_gateway(self, payload: Any) -> dict[str, Any]:
  276. body = _closed(payload, {"subject_uid", "roles", "business_domain_uid", "asset_uid", "purpose", "environment", "action", "classification", "requested_fields"}, "gateway request")
  277. context = {
  278. "subject_uid": _uid(body.get("subject_uid"), "subject_uid"),
  279. "roles": sorted({_text(item, "role", 80) for item in _values(body.get("roles"), "roles", minimum=0)}),
  280. "business_domain_uid": _uid(body.get("business_domain_uid"), "business_domain_uid"),
  281. "asset_uid": _uid(body.get("asset_uid"), "asset_uid"), "purpose": _text(body.get("purpose"), "purpose", 100),
  282. "environment": _text(body.get("environment"), "environment", 30), "action": _text(body.get("action"), "action", 20),
  283. "classification": _text(body.get("classification"), "classification", 30),
  284. "requested_fields": sorted({_field(item, "requested field") for item in _values(body.get("requested_fields"), "requested fields", maximum=500)}),
  285. }
  286. if context["environment"] not in _ENVIRONMENTS or context["action"] not in _ACTIONS or context["classification"] not in _CLASSIFICATIONS:
  287. raise ValueError("unsupported gateway environment, action, or classification")
  288. if context["action"] == "export" and self.repository.active_hold_for_asset(context["asset_uid"]):
  289. return {"decision": "deny", "reason_code": "active_legal_hold"}
  290. if context["action"] == "export" and context["classification"] == "highly_sensitive":
  291. return {"decision": "deny", "reason_code": "highly_sensitive_export_denied"}
  292. now = self.now_factory().astimezone(UTC)
  293. for policy in self.repository.active_policies():
  294. selector = policy["selector"]
  295. if _time(selector["expires_at"], "policy expiry") <= now:
  296. continue
  297. if not (context["subject_uid"] in selector["subjects"] or set(context["roles"]) & set(selector["roles"])):
  298. continue
  299. selector_keys = {
  300. "business_domain_uid": "business_domains", "asset_uid": "assets",
  301. "purpose": "purposes", "environment": "environments",
  302. "action": "actions", "classification": "classifications",
  303. }
  304. if any(context[key] not in selector[selector_keys[key]] for key in selector_keys):
  305. continue
  306. rules = {rule["field"]: rule for rule in policy["resource_rules"]["field_rules"]}
  307. if any(field not in rules or rules[field]["effect"] != "allow" for field in context["requested_fields"]):
  308. continue
  309. masks = {field: rules[field]["mask_ref"] for field in context["requested_fields"] if rules[field]["mask_ref"]}
  310. return {"decision": "allow", "reason_code": "policy_allowed", "policy_code": policy["code"], "policy_version": policy["version"], "projected_fields": context["requested_fields"], "row_filter": policy["resource_rules"]["row_filter"], "field_masks": masks}
  311. return {"decision": "deny", "reason_code": "default_deny"}
  312. def create_grant(self, payload: Any, *, actor_uid: str) -> dict[str, Any]:
  313. body = _closed(payload, {"policy_uid", "approval_ref", "approval_digest", "subject_uid", "asset_uid", "provider", "expires_at", "decision"}, "trusted delivery grant")
  314. actor = self._actor(actor_uid)
  315. decision = _closed(body.get("decision"), {"decision", "reason_code", "policy_code", "policy_version", "projected_fields", "row_filter", "field_masks"}, "grant decision")
  316. if decision.get("decision") != "allow":
  317. raise PermissionError("trusted delivery grant requires allowed decision")
  318. provider = _text(body.get("provider"), "provider", 30)
  319. self.provider_registry.get(provider)
  320. expires_at = _time(body.get("expires_at"), "grant expiry")
  321. if expires_at <= self.now_factory().astimezone(UTC):
  322. raise ValueError("grant expiry must be in the future")
  323. record = {
  324. "uid": self.uid_factory(), "policy_uid": _uid(body.get("policy_uid"), "policy_uid"),
  325. "approval_ref": _text(body.get("approval_ref"), "approval reference", 200),
  326. "approval_digest": _digest_value(body.get("approval_digest"), "approval digest"),
  327. "subject_uid": _uid(body.get("subject_uid"), "subject_uid"), "asset_uid": _uid(body.get("asset_uid"), "asset_uid"),
  328. "provider": provider, "expires_at": expires_at.isoformat(), "status": "active", "reason_code": "approval_bound",
  329. "decision_digest": _digest(decision), "current_version": 1, "created_by": actor, "created_at": self.now_factory().isoformat(),
  330. }
  331. saved = self.repository.create_grant(record)
  332. result = self.repository.get_grant(saved["uid"]) if hasattr(self.repository, "get_grant") else saved
  333. self._event("grant_created", actor, {"grant_uid": result["uid"], "provider": provider, "decision_digest": record["decision_digest"]})
  334. return result
  335. def dispatch_grant(self, grant_uid: str, *, actor_uid: str) -> dict[str, Any]:
  336. actor = self._actor(actor_uid)
  337. if hasattr(self.repository, "get_grant"):
  338. grant = self.repository.get_grant(_uid(grant_uid, "grant_uid"))
  339. elif hasattr(self.repository, "grants"):
  340. grant = copy.deepcopy(self.repository.grants.get(_uid(grant_uid, "grant_uid")))
  341. else:
  342. raise RuntimeError("durable dispatch repository is required")
  343. if not grant or grant["status"] != "active":
  344. raise LookupError("active trusted delivery grant was not found")
  345. response = self.provider_registry.get(grant["provider"]).apply({
  346. "schema": "dataops.trusted-delivery.v1", "grant_uid": grant["uid"], "provider": grant["provider"],
  347. "subject_uid": grant["subject_uid"], "asset_uid": grant["asset_uid"], "expires_at": grant["expires_at"], "decision_digest": grant["decision_digest"],
  348. })
  349. response = _closed(response, {"status", "receipt_code", "response_digest"}, "provider receipt")
  350. if response.get("status") != "applied":
  351. raise RuntimeError("trusted delivery provider failed closed")
  352. receipt = {"grant_uid": grant["uid"], "status": "applied", "receipt_code": _text(response.get("receipt_code"), "receipt code", 120), "response_digest": _digest_value(response.get("response_digest"), "provider response digest")}
  353. self._event("grant_dispatched", actor, receipt)
  354. return receipt
  355. def provision(self, payload: Any, *, actor_uid: str, subject_roles: list[str] | None = None) -> dict[str, Any]:
  356. """Create and apply one approval-bound, digest-only capability grant."""
  357. body = _closed(payload, {"policy_uid", "approval_ref", "approval_digest", "subject_uid", "asset_uid", "business_domain_uid", "purpose", "environment", "action", "classification", "requested_fields", "expires_at", "target_capability", "provider", "idempotency_key"}, "trusted delivery provision")
  358. actor = self._actor(actor_uid)
  359. provider = _text(body.get("provider"), "provider", 30)
  360. adapter = self.provider_registry.get(provider)
  361. capability = _closed(body.get("target_capability"), {"actions", "fields"}, "target capability")
  362. target_capability = {
  363. "actions": _simple_values(capability.get("actions"), "capability actions", _ACTIONS),
  364. "fields": sorted({_field(item, "capability field") for item in _values(capability.get("fields"), "capability fields", maximum=500)}),
  365. }
  366. approval_ref = _text(body.get("approval_ref"), "approval reference", 200)
  367. approval_digest = _digest_value(body.get("approval_digest"), "approval digest")
  368. policy_uid = _uid(body.get("policy_uid"), "policy_uid")
  369. subject_uid = _uid(body.get("subject_uid"), "subject_uid")
  370. asset_uid = _uid(body.get("asset_uid"), "asset_uid")
  371. context = {
  372. "subject_uid": subject_uid,
  373. "roles": list(subject_roles or []),
  374. "business_domain_uid": _uid(body.get("business_domain_uid"), "business_domain_uid"),
  375. "asset_uid": asset_uid,
  376. "purpose": _text(body.get("purpose"), "purpose", 100),
  377. "environment": _text(body.get("environment"), "environment", 30),
  378. "action": _text(body.get("action"), "action", 20),
  379. "classification": _text(body.get("classification"), "classification", 30),
  380. "requested_fields": sorted({_field(item, "requested field") for item in _values(body.get("requested_fields"), "requested fields", maximum=500)}),
  381. }
  382. # Re-evaluate from the durable, active policy version. A caller's
  383. # policy UID, capability and approval labels are only references; none
  384. # of them is allowed to enlarge the authorisation projection.
  385. decision = self.evaluate_gateway(context)
  386. active = next((item for item in self.repository.active_policies() if item.get("uid") == policy_uid), None)
  387. if active is None or decision.get("decision") != "allow" or active["code"] != decision.get("policy_code") or active["version"] != decision.get("policy_version"):
  388. raise PermissionError("trusted delivery provision is not allowed by active policy")
  389. if set(target_capability["actions"]) != {context["action"]} or set(target_capability["fields"]) != set(decision["projected_fields"]):
  390. raise PermissionError("target capability exceeds server policy projection")
  391. expires_at = _time(body.get("expires_at"), "grant expiry")
  392. if expires_at <= self.now_factory().astimezone(UTC):
  393. raise ValueError("grant expiry must be in the future")
  394. policy_expiry = _time(active["selector"]["expires_at"], "policy expiry")
  395. if expires_at > policy_expiry:
  396. raise PermissionError("grant expiry exceeds active policy")
  397. if "export" in target_capability["actions"] and self.repository.active_hold_for_asset(asset_uid):
  398. raise PermissionError("active legal hold blocks export provision")
  399. projection = {
  400. "fields": decision["projected_fields"],
  401. "row_predicate": decision["row_filter"],
  402. "masking": decision["field_masks"],
  403. }
  404. decision_digest = _digest({"policy_uid": policy_uid, "policy_code": active["code"], "policy_version": active["version"], "context": context, "projection": projection, "approval_digest": approval_digest})
  405. record = {
  406. "uid": self.uid_factory(), "policy_uid": policy_uid,
  407. "approval_ref": approval_ref, "approval_digest": approval_digest,
  408. "subject_uid": subject_uid, "asset_uid": asset_uid,
  409. "purpose": context["purpose"], "environment": context["environment"],
  410. "target_capability": target_capability, "provider": provider,
  411. "idempotency_key": _text(body.get("idempotency_key"), "idempotency key", 160),
  412. "expires_at": expires_at.isoformat(), "status": "active", "reason_code": "approval_bound",
  413. "request_digest": _digest({"decision_digest": decision_digest, "provider": provider, "expires_at": expires_at.isoformat(), "idempotency_key": body.get("idempotency_key")}),
  414. "decision_digest": decision_digest,
  415. "current_version": 1, "created_by": actor,
  416. }
  417. if hasattr(self.repository, "enqueue_grant"):
  418. grant = self.repository.enqueue_grant(record)
  419. grant_uid = grant["uid"]
  420. if hasattr(self.repository, "get_grant"):
  421. stored = self.repository.get_grant(grant_uid)
  422. if stored and stored.get("status") == "reclaimed":
  423. raise RuntimeError("trusted delivery grant was reclaimed")
  424. if hasattr(self.repository, "get_provision_receipt"):
  425. prior_receipt = self.repository.get_provision_receipt(grant_uid, record["idempotency_key"])
  426. if prior_receipt:
  427. return {"grant_uid": grant_uid, "status": "applied", "receipt": prior_receipt}
  428. else:
  429. raise RuntimeError("durable provision repository is required")
  430. # Persist a fenced, idempotent outbox intent before the adapter is
  431. # permitted to see the request. A process crash now leaves a retryable
  432. # pending record rather than an unprovable external side effect.
  433. if hasattr(self.repository, "prepare_provision_operation"):
  434. operation = self.repository.prepare_provision_operation(
  435. grant_uid, record["idempotency_key"], record["request_digest"]
  436. )
  437. if operation.get("completed"):
  438. return {"grant_uid": grant_uid, "status": "applied", "receipt": operation["receipt"]}
  439. elif hasattr(self.repository, "session"):
  440. raise RuntimeError("durable provision outbox repository is required")
  441. envelope = {"schema": "dataops.trusted-delivery.v2", "operation": "provision", "grant_uid": grant_uid, "policy_uid": record["policy_uid"], "provider": provider, "subject_uid": record["subject_uid"], "asset_uid": record["asset_uid"], "purpose": record["purpose"], "environment": record["environment"], "expires_at": record["expires_at"], "projection": projection, "decision_digest": decision_digest, "request_digest": record["request_digest"]}
  442. claim = None
  443. if hasattr(self.repository, "claim_provision_operation"):
  444. worker = f"provision:{self.uid_factory()}"
  445. claim = self.repository.claim_provision_operation(operation["delivery_uid"], worker)
  446. if claim.get("state") == "completed":
  447. return {"grant_uid": grant_uid, "status": "applied", "receipt": claim["receipt"]}
  448. if claim.get("state") != "claimed":
  449. # Another process owns the external side effect. Returning a
  450. # pending state is intentionally safer than performing a
  451. # duplicate call while that DB-time lease is alive.
  452. return {"grant_uid": grant_uid, "status": "pending", "retryable": True}
  453. try:
  454. response = _closed(adapter.apply(copy.deepcopy(envelope)), {"status", "receipt_code", "response_digest"}, "provider receipt")
  455. if response.get("status") != "applied":
  456. raise RuntimeError("trusted delivery provider failed closed")
  457. except Exception:
  458. if claim is not None:
  459. # Commit the retryable outbox state before the HTTP boundary
  460. # rolls back the failed request transaction.
  461. self.repository.fail_provision_operation(operation["delivery_uid"], worker, claim["lease_fence"])
  462. self.repository.session.commit()
  463. raise
  464. receipt = {"grant_uid": grant_uid, "status": "applied", "receipt_code": _text(response.get("receipt_code"), "receipt code", 120), "response_digest": _digest_value(response.get("response_digest"), "provider response digest"), "diff_digest": _digest({"operation": "provision", "request_digest": record["request_digest"], "response_digest": response.get("response_digest")})}
  465. if claim is not None:
  466. saved = self.repository.complete_provision_operation(
  467. grant_uid, operation["delivery_uid"], worker, claim["lease_fence"], receipt
  468. )
  469. receipt = {**receipt, "grant_uid": saved.get("uid", grant_uid)}
  470. elif hasattr(self.repository, "record_provision_receipt"):
  471. receipt = self.repository.record_provision_receipt(grant_uid, record["idempotency_key"], record["request_digest"], receipt)
  472. self._event("provision_applied", actor, {"grant_uid": grant_uid, "provider": provider, "request_digest": record["request_digest"], "response_digest": receipt["response_digest"]})
  473. return {"grant_uid": grant_uid, "status": "applied", "receipt": receipt}
  474. def create_legal_hold(self, payload: Any, *, actor_uid: str) -> dict[str, Any]:
  475. body = _closed(payload, {"asset_uid", "operation", "approver_refs", "evidence_digest"}, "legal hold")
  476. actor = self._actor(actor_uid)
  477. asset_uid = _uid(body.get("asset_uid"), "asset_uid")
  478. if self.repository.active_hold_for_asset(asset_uid):
  479. raise RuntimeError("active legal hold already exists")
  480. operation = _text(body.get("operation"), "legal hold operation", 30)
  481. if operation not in {"freeze", "export_review", "destruction_review"}:
  482. raise ValueError("unsupported legal hold operation")
  483. approvers = sorted({_text(item, "approver reference", 200) for item in _values(body.get("approver_refs"), "approver references", minimum=2, maximum=2)})
  484. if len(approvers) != 2:
  485. raise ValueError("legal hold requires two distinct approver references")
  486. record = {
  487. "uid": self.uid_factory(), "asset_uid": asset_uid, "operation": operation,
  488. # The SQL schema intentionally stores the two approvals in separate
  489. # columns so a parameter mapping cannot silently collapse them.
  490. "approver_ref_one": approvers[0], "approver_ref_two": approvers[1],
  491. "approver_refs": approvers, "evidence_digest": _digest_value(body.get("evidence_digest"), "evidence digest"),
  492. "status": "active", "created_by": actor, "created_at": self.now_factory().isoformat(), "current_version": 1,
  493. }
  494. saved = self.repository.create_legal_hold(record)
  495. self._event("legal_hold_active", actor, {"hold_uid": saved["uid"], "operation": operation})
  496. return saved
  497. def release_legal_hold(self, payload: Any, *, actor_uid: str) -> dict[str, Any]:
  498. body = _closed(payload, {"hold_uid", "approval_ref", "approval_digest", "idempotency_key"}, "legal hold release")
  499. actor = self._actor(actor_uid)
  500. hold_uid = _uid(body.get("hold_uid"), "hold_uid")
  501. approval_ref = _text(body.get("approval_ref"), "approval reference", 200)
  502. approval_digest = _digest_value(body.get("approval_digest"), "approval digest")
  503. key = _text(body.get("idempotency_key"), "idempotency key", 160)
  504. request_digest = _digest({"hold_uid": hold_uid, "approval_ref": approval_ref, "approval_digest": approval_digest})
  505. if hasattr(self.repository, "release_legal_hold"):
  506. saved = self.repository.release_legal_hold(hold_uid, key, request_digest, actor, approval_ref, approval_digest)
  507. elif hasattr(self.repository, "holds"):
  508. releases = getattr(self.repository, "hold_releases", {})
  509. existing = releases.get(key)
  510. if existing:
  511. if existing["request_digest"] != request_digest:
  512. raise RuntimeError("trusted delivery idempotency conflict")
  513. return copy.deepcopy(existing["result"])
  514. hold = self.repository.holds.get(hold_uid)
  515. if not hold:
  516. raise LookupError("trusted delivery legal hold was not found")
  517. if hold["created_by"] == actor or approval_ref in hold["approver_refs"]:
  518. raise PermissionError("legal hold release requires an independent approval")
  519. hold["status"] = "released"
  520. hold["current_version"] += 1
  521. saved = {"uid": hold_uid, "status": "released", "current_version": hold["current_version"]}
  522. releases[key] = {"request_digest": request_digest, "result": copy.deepcopy(saved)}
  523. self.repository.hold_releases = releases
  524. else:
  525. raise RuntimeError("durable legal hold repository is required")
  526. self._event("legal_hold_released", actor, {"hold_uid": hold_uid, "approval_digest": approval_digest})
  527. return saved
  528. def execute_reclaim(self, payload: Any, *, actor_uid: str) -> dict[str, Any]:
  529. """Request provider revocation first; never delete or reclaim optimistically."""
  530. body = _closed(payload, {"grant_uid", "idempotency_key"}, "trusted delivery reclaim")
  531. actor = self._actor(actor_uid)
  532. grant_uid = _uid(body.get("grant_uid"), "grant_uid")
  533. key = _text(body.get("idempotency_key"), "idempotency key", 160)
  534. if not hasattr(self.repository, "get_grant"):
  535. raise RuntimeError("durable reclaim repository is required")
  536. grant = self.repository.get_grant(grant_uid)
  537. if not grant:
  538. raise LookupError("trusted delivery grant was not found")
  539. now = self.now_factory().astimezone(UTC)
  540. if grant.get("status") == "reclaimed":
  541. return {"grant_uid": grant_uid, "status": "reclaimed", "idempotent": True}
  542. if grant.get("status") != "active" or _time(grant["expires_at"], "grant expiry") > now:
  543. raise ValueError("trusted delivery grant is not reclaimable")
  544. if self.repository.active_hold_for_asset(grant["asset_uid"]):
  545. raise PermissionError("active legal hold blocks reclaim")
  546. expiry = _time(grant["expires_at"], "grant expiry").isoformat()
  547. request_digest = _digest({"operation": "revoke", "grant_uid": grant_uid, "provider": grant["provider"], "expires_at": expiry})
  548. worker_id = f"reclaim-{uuid.uuid4()}"
  549. claim = None
  550. if hasattr(self.repository, "claim_reclaim_operation"):
  551. claim = self.repository.claim_reclaim_operation(grant_uid, key, request_digest, worker_id)
  552. if claim["state"] == "completed":
  553. return {"grant_uid": grant_uid, "status": "reclaimed", "idempotent": True}
  554. if claim["state"] == "leased":
  555. raise RuntimeError("trusted delivery reclaim is already leased")
  556. try:
  557. adapter = self.provider_registry.get(grant["provider"])
  558. except PermissionError:
  559. self._event("reclaim_provider_unavailable", actor, {"grant_uid": grant_uid, "provider": grant["provider"], "request_digest": request_digest})
  560. raise
  561. envelope = {"schema": "dataops.trusted-delivery.v2", "operation": "revoke", "grant_uid": grant_uid, "provider": grant["provider"], "asset_uid": grant["asset_uid"], "request_digest": request_digest}
  562. revoke = getattr(adapter, "revoke", None)
  563. if not callable(revoke):
  564. self._event("reclaim_provider_unavailable", actor, {"grant_uid": grant_uid, "provider": grant["provider"], "request_digest": request_digest})
  565. raise PermissionError("trusted delivery provider revoke is disabled")
  566. try:
  567. response = _closed(revoke(copy.deepcopy(envelope)), {"status", "receipt_code", "response_digest"}, "provider receipt")
  568. except Exception:
  569. if claim:
  570. self.repository.fail_reclaim_operation(claim["delivery_uid"], worker_id, claim["lease_fence"])
  571. raise
  572. if response.get("status") != "revoked":
  573. if claim:
  574. self.repository.fail_reclaim_operation(claim["delivery_uid"], worker_id, claim["lease_fence"])
  575. raise RuntimeError("trusted delivery provider failed closed")
  576. receipt = {"grant_uid": grant_uid, "status": "reclaimed", "receipt_code": _text(response.get("receipt_code"), "receipt code", 120), "response_digest": _digest_value(response.get("response_digest"), "provider response digest"), "diff_digest": _digest({"operation": "revoke", "request_digest": request_digest, "response_digest": response.get("response_digest")})}
  577. if claim:
  578. saved = self.repository.complete_claimed_reclaim(grant_uid, claim["delivery_uid"], worker_id, claim["lease_fence"], receipt)
  579. elif hasattr(self.repository, "complete_reclaim"):
  580. saved = self.repository.complete_reclaim(grant_uid, key, request_digest, receipt)
  581. else:
  582. saved = self.repository.update_grant_status(grant_uid, grant["current_version"], "reclaimed", "expired_reclaimed")
  583. self._event("grant_reclaimed", actor, {"grant_uid": grant_uid, "provider": grant["provider"], "request_digest": request_digest, "response_digest": receipt["response_digest"]})
  584. return {"grant_uid": grant_uid, "status": saved["status"], "receipt": receipt}
  585. def preview_reclaim(self, *, as_of: Any, actor_uid: str) -> list[dict[str, Any]]:
  586. """Return eligible identifiers only; this query has no state transition."""
  587. self._actor(actor_uid)
  588. instant = _time(as_of, "as_of")
  589. candidates = []
  590. for grant in self.repository.grants_due_for_revoke(instant):
  591. if grant.get("status") != "active" or _time(grant["expires_at"], "grant expiry") > instant:
  592. continue
  593. if self.repository.active_hold_for_asset(grant["asset_uid"]):
  594. continue
  595. candidates.append({"grant_uid": grant["uid"], "provider": grant.get("provider", "unknown"), "status": "eligible"})
  596. return candidates
  597. def revoke_expired(self, *, as_of: datetime, actor_uid: str) -> list[dict[str, Any]]:
  598. actor = self._actor(actor_uid)
  599. instant = _time(as_of, "as_of")
  600. revoked = []
  601. for grant in self.repository.grants_due_for_revoke(instant):
  602. if _time(grant["expires_at"], "grant expiry") > instant or self.repository.active_hold_for_asset(grant["asset_uid"]):
  603. continue
  604. saved = self.repository.update_grant_status(grant["uid"], grant["current_version"], "revoked", "expired")
  605. self._event("grant_revoked", actor, {"grant_uid": saved["uid"], "reason_code": "expired"})
  606. revoked.append(saved)
  607. return revoked
  608. @staticmethod
  609. def mask_preview(payload: Any) -> dict[str, Any]:
  610. body = _closed(payload, {"mode", "mask_ref", "value"}, "mask preview")
  611. mode = _text(body.get("mode"), "mask mode", 20)
  612. if mode not in {"static", "dynamic", "display", "export"}:
  613. raise ValueError("unsupported mask mode")
  614. mask_ref = _text(body.get("mask_ref"), "mask ref", 120)
  615. value = _text(body.get("value"), "mask value", 1000)
  616. return {"mode": mode, "mask_ref": mask_ref, "masked_value": "***-****", "token_digest": _digest({"mask_ref": mask_ref, "value": value})}