"""Closed, default-deny contracts for fine-grained trusted delivery. This module intentionally never executes a data query or contacts a target. It produces a bounded authorization projection and digest-only provider envelope. Concrete enterprise providers remain external integration work. """ from __future__ import annotations import copy import hashlib import json import re import uuid from collections.abc import Callable, Mapping from datetime import UTC, datetime from typing import Any from app.core.common.identifiers import new_governance_uid from app.core.common.timezone_utils import now_china _PROVIDERS = frozenset({"database", "api", "file", "iam", "kms", "dlp", "siem"}) _ENVIRONMENTS = frozenset({"development", "test", "production"}) _ACTIONS = frozenset({"read", "use", "display", "export"}) _CLASSIFICATIONS = frozenset({"public", "internal", "sensitive", "highly_sensitive"}) _FIELD = re.compile(r"^[A-Za-z_][A-Za-z0-9_.-]{0,199}$") _CODE = re.compile(r"^[A-Z][A-Z0-9_]{2,119}$") _DIGEST = re.compile(r"^[0-9a-f]{64}$") _UNSAFE_TEXT = re.compile( r"(?i)(?:https?://|ftp://|file://|www\.|\b(?:secret|token|password|private[ _-]?key)\b|" r"-----BEGIN| dict[str, Any]: if not isinstance(value, dict): raise ValueError(f"{label} must be an object") extra = sorted(set(value) - allowed) if extra: raise ValueError(f"{label} contains unsupported fields: {', '.join(extra)}") result = copy.deepcopy(value) _assert_safe_value(result, label) return result def _assert_safe_value(value: Any, label: str) -> None: """Reject values that could turn an authorization reference into payload data.""" if isinstance(value, str) and _UNSAFE_TEXT.search(value): raise ValueError(f"{label} contains unsafe content") if isinstance(value, dict): for item in value.values(): _assert_safe_value(item, label) elif isinstance(value, list): for item in value: _assert_safe_value(item, label) def _text(value: Any, label: str, maximum: int = 300) -> str: if not isinstance(value, str) or not value.strip(): raise ValueError(f"{label} is required") result = value.strip() if len(result) > maximum: raise ValueError(f"{label} exceeds {maximum} characters") return result def _uid(value: Any, label: str) -> str: try: return str(uuid.UUID(str(value))) except (TypeError, ValueError, AttributeError) as exc: raise ValueError(f"{label} must be a UUID") from exc def _time(value: Any, label: str) -> datetime: try: result = value if isinstance(value, datetime) else datetime.fromisoformat(str(value).replace("Z", "+00:00")) except (TypeError, ValueError) as exc: raise ValueError(f"{label} must be ISO-8601") from exc if result.tzinfo is None: raise ValueError(f"{label} must include a timezone") return result.astimezone(UTC) def _values(value: Any, label: str, *, minimum: int = 1, maximum: int = 200) -> list[Any]: if not isinstance(value, list) or not minimum <= len(value) <= maximum: raise ValueError(f"{label} must contain between {minimum} and {maximum} items") return copy.deepcopy(value) def _digest(value: Any) -> str: return hashlib.sha256(json.dumps(value, sort_keys=True, separators=(",", ":")).encode()).hexdigest() def _digest_value(value: Any, label: str) -> str: result = _text(value, label, 64).lower() if not _DIGEST.fullmatch(result): raise ValueError(f"{label} must be SHA-256") return result def _field(value: Any, label: str) -> str: result = _text(value, label, 200) if not _FIELD.fullmatch(result): raise ValueError(f"{label} is invalid") return result def _simple_values(value: Any, label: str, allowed: frozenset[str]) -> list[str]: result = sorted({_text(item, label, 100) for item in _values(value, label)}) if not set(result) <= allowed: raise ValueError(f"unsupported {label}") return result def _row_filter(value: Any) -> dict[str, Any] | None: if value is None: return None body = _closed(value, {"op", "field", "value", "items", "item"}, "row filter") op = _text(body.get("op"), "row filter op", 10) if op not in {"eq", "in", "and", "or", "not"}: raise ValueError("unsupported row filter op") if _UNSAFE_SQL.search(json.dumps(body, sort_keys=True)): raise ValueError("row filter contains unsafe content") if op in {"eq", "in"}: if set(body) - {"op", "field", "value"}: raise ValueError("row filter shape is invalid") field = _field(body.get("field"), "row filter field") raw = body.get("value") if op == "eq" and not isinstance(raw, (str, int, float, bool)): raise ValueError("row filter eq value must be scalar") if op == "in": if not isinstance(raw, list) or not 1 <= len(raw) <= 100 or not all(isinstance(item, (str, int, float, bool)) for item in raw): raise ValueError("row filter in value must be bounded scalar list") raw = sorted(raw, key=lambda item: (type(item).__name__, str(item))) return {"op": op, "field": field, "value": raw} if op == "not": if set(body) != {"op", "item"}: raise ValueError("row filter not shape is invalid") return {"op": op, "item": _row_filter(body["item"])} if set(body) != {"op", "items"}: raise ValueError("row filter boolean shape is invalid") items = [_row_filter(item) for item in _values(body.get("items"), "row filter items", minimum=1, maximum=20)] return {"op": op, "items": items} class ClosedProviderRegistry: """Provider names are a closed enum; tests may inject non-network adapters.""" def __init__(self, adapters: Mapping[str, Any] | None = None, *, test_mode: bool = False): adapters = dict(adapters or {}) if set(adapters) - _PROVIDERS: raise ValueError("unsupported trusted-delivery provider") if adapters and not test_mode: raise PermissionError("enterprise providers require approved external integration") self._adapters = adapters self._test_mode = test_mode @classmethod def for_tests(cls, adapters: Mapping[str, Any]): return cls(adapters, test_mode=True) @property def is_enterprise_ready(self) -> bool: return False def dependency_status(self) -> dict[str, str]: return dict.fromkeys(sorted(_PROVIDERS), "TBD_EXTERNAL") def get(self, provider: str): provider = _text(provider, "provider", 30) if provider not in _PROVIDERS: raise ValueError("unsupported trusted-delivery provider") adapter = self._adapters.get(provider) if adapter is None: raise PermissionError(f"trusted-delivery provider {provider} is disabled") return adapter class TrustedDeliveryService: """Durable-service façade. Repositories own the database transaction boundary.""" def __init__( self, repository, *, provider_registry: ClosedProviderRegistry | None = None, uid_factory: Callable[[], str] = new_governance_uid, now_factory: Callable[[], datetime] = now_china, ): self.repository = repository self.provider_registry = provider_registry or ClosedProviderRegistry() self.uid_factory = uid_factory self.now_factory = now_factory def _actor(self, value: Any) -> str: actor = _uid(value, "actor_uid") if self.repository.users_available({actor}) != {actor}: raise ValueError("trusted-delivery actor is unavailable") return actor def _event(self, action: str, actor_uid: str, detail: dict[str, Any]) -> None: if hasattr(self.repository, "add_evidence"): self.repository.add_evidence({ "uid": self.uid_factory(), "action": action, "actor_uid": actor_uid, "safe_detail": copy.deepcopy(detail), "created_at": self.now_factory().isoformat(), }) def create_policy_version(self, payload: Any, *, actor_uid: str) -> dict[str, Any]: body = _closed(payload, {"code", "version", "selector", "resource_rules"}, "trusted delivery policy") actor = self._actor(actor_uid) code = _text(body.get("code"), "policy code", 120).upper() if not _CODE.fullmatch(code): raise ValueError("policy code is invalid") selector = _closed(body.get("selector"), {"subjects", "roles", "business_domains", "assets", "purposes", "environments", "actions", "classifications", "expires_at"}, "policy selector") subjects = sorted({_uid(item, "policy subject") for item in _values(selector.get("subjects", []), "policy subjects", minimum=0)}) roles = sorted({_text(item, "policy role", 80) for item in _values(selector.get("roles", []), "policy roles", minimum=0)}) if not subjects and not roles: raise ValueError("policy selector requires a subject or role") if subjects and self.repository.users_available(set(subjects)) != set(subjects): raise ValueError("policy selector contains unavailable users") rules = _closed(body.get("resource_rules"), {"row_filter", "field_rules"}, "resource rules") field_rules = [] for item in _values(rules.get("field_rules"), "field rules", maximum=500): field_rule = _closed(item, {"field", "effect", "mask_ref"}, "field rule") effect = _text(field_rule.get("effect"), "field effect", 10) if effect not in {"allow", "deny"}: raise ValueError("unsupported field effect") mask_ref = field_rule.get("mask_ref") if mask_ref is not None: mask_ref = _text(mask_ref, "mask ref", 120) if not re.fullmatch(r"[a-z][a-z0-9-]{2,119}-v[1-9][0-9]*", mask_ref): raise ValueError("mask ref is invalid") field_rules.append({"field": _field(field_rule.get("field"), "field rule field"), "effect": effect, "mask_ref": mask_ref}) if len({item["field"] for item in field_rules}) != len(field_rules): raise ValueError("field rules must not duplicate fields") expires_at = _time(selector.get("expires_at"), "policy expiry") now = self.now_factory().astimezone(UTC) if expires_at <= now: raise ValueError("policy expiry must be in the future") record = { "uid": self.uid_factory(), "code": code, "version": _text(body.get("version"), "policy version", 30), # A policy never becomes effective by creation alone. This keeps # both first and subsequent versions behind the durable approval # transition, rather than treating an empty policy history as an # implicit approval. "status": "draft", "selector": { "subjects": subjects, "roles": roles, "business_domains": sorted({_uid(item, "business domain") for item in _values(selector.get("business_domains"), "business domains")}), "assets": sorted({_uid(item, "asset") for item in _values(selector.get("assets"), "assets")}), "purposes": sorted({_text(item, "purpose", 100) for item in _values(selector.get("purposes"), "purposes")}), "environments": _simple_values(selector.get("environments"), "environments", _ENVIRONMENTS), "actions": _simple_values(selector.get("actions"), "actions", _ACTIONS), "classifications": _simple_values(selector.get("classifications"), "classifications", _CLASSIFICATIONS), "expires_at": expires_at.isoformat(), }, "resource_rules": {"row_filter": _row_filter(rules.get("row_filter")), "field_rules": sorted(field_rules, key=lambda item: item["field"])}, "created_by": actor, "created_at": now.isoformat(), "current_version": 1, } saved = self.repository.create_policy_version(record) self._event("policy_version_created", actor, {"policy_code": code, "version": record["version"]}) return saved def _lifecycle_request(self, payload: Any, *, label: str) -> tuple[dict[str, Any], str, str, str]: body = _closed(payload, {"code", "version", "approval_ref", "approval_digest", "idempotency_key"}, label) actor = self._actor(payload.get("actor_uid")) if isinstance(payload, dict) and "actor_uid" in payload else None # Callers supply the actor separately. This branch only keeps direct # service use closed if a malformed payload is passed. if actor is not None: raise ValueError("actor_uid must not be in lifecycle payload") return ( body, _text(body.get("code"), "policy code", 120).upper(), _text(body.get("version"), "policy version", 30), _text(body.get("idempotency_key"), "idempotency key", 160), ) def _policy_transition(self, payload: Any, *, actor_uid: str, operation: str) -> dict[str, Any]: body = _closed(payload, {"code", "version", "approval_ref", "approval_digest", "idempotency_key"}, "policy transition") actor = self._actor(actor_uid) code = _text(body.get("code"), "policy code", 120).upper() if not _CODE.fullmatch(code): raise ValueError("policy code is invalid") version = _text(body.get("version"), "policy version", 30) approval_ref = _text(body.get("approval_ref"), "approval reference", 200) approval_digest = _digest_value(body.get("approval_digest"), "approval digest") key = _text(body.get("idempotency_key"), "idempotency key", 160) request_digest = _digest({"operation": operation, "code": code, "version": version, "approval_ref": approval_ref, "approval_digest": approval_digest}) if hasattr(self.repository, "transition_policy"): saved = self.repository.transition_policy(code, version, operation, key, request_digest, actor, approval_ref, approval_digest) elif hasattr(self.repository, "policies"): # Small test repositories intentionally do not emulate PostgreSQL; # preserve the same state rules for service-level tests. history = getattr(self.repository, "policy_transitions", {}) existing = history.get(key) if existing: if existing["request_digest"] != request_digest: raise RuntimeError("trusted delivery idempotency conflict") return copy.deepcopy(existing["result"]) target = next((item for item in self.repository.policies.values() if item["code"] == code and item["version"] == version), None) if target is None: raise LookupError("trusted delivery policy version was not found") for item in self.repository.policies.values(): if item["code"] == code and item["uid"] != target["uid"] and item["status"] == "active": item["status"] = "retired" item["current_version"] += 1 if target["status"] != "active": target["status"] = "active" target["current_version"] += 1 saved = {key: target[key] for key in ("uid", "code", "version", "status", "current_version")} history[key] = {"request_digest": request_digest, "result": copy.deepcopy(saved)} self.repository.policy_transitions = history else: raise RuntimeError("durable policy transition repository is required") self._event(f"policy_{operation}", actor, {"policy_code": code, "version": version, "approval_digest": approval_digest}) return saved def activate_policy_version(self, payload: Any, *, actor_uid: str) -> dict[str, Any]: return self._policy_transition(payload, actor_uid=actor_uid, operation="activate") def rollback_policy_version(self, payload: Any, *, actor_uid: str) -> dict[str, Any]: return self._policy_transition(payload, actor_uid=actor_uid, operation="rollback") def evaluate_gateway(self, payload: Any) -> dict[str, Any]: body = _closed(payload, {"subject_uid", "roles", "business_domain_uid", "asset_uid", "purpose", "environment", "action", "classification", "requested_fields"}, "gateway request") context = { "subject_uid": _uid(body.get("subject_uid"), "subject_uid"), "roles": sorted({_text(item, "role", 80) for item in _values(body.get("roles"), "roles", minimum=0)}), "business_domain_uid": _uid(body.get("business_domain_uid"), "business_domain_uid"), "asset_uid": _uid(body.get("asset_uid"), "asset_uid"), "purpose": _text(body.get("purpose"), "purpose", 100), "environment": _text(body.get("environment"), "environment", 30), "action": _text(body.get("action"), "action", 20), "classification": _text(body.get("classification"), "classification", 30), "requested_fields": sorted({_field(item, "requested field") for item in _values(body.get("requested_fields"), "requested fields", maximum=500)}), } if context["environment"] not in _ENVIRONMENTS or context["action"] not in _ACTIONS or context["classification"] not in _CLASSIFICATIONS: raise ValueError("unsupported gateway environment, action, or classification") if context["action"] == "export" and self.repository.active_hold_for_asset(context["asset_uid"]): return {"decision": "deny", "reason_code": "active_legal_hold"} if context["action"] == "export" and context["classification"] == "highly_sensitive": return {"decision": "deny", "reason_code": "highly_sensitive_export_denied"} now = self.now_factory().astimezone(UTC) for policy in self.repository.active_policies(): selector = policy["selector"] if _time(selector["expires_at"], "policy expiry") <= now: continue if not (context["subject_uid"] in selector["subjects"] or set(context["roles"]) & set(selector["roles"])): continue selector_keys = { "business_domain_uid": "business_domains", "asset_uid": "assets", "purpose": "purposes", "environment": "environments", "action": "actions", "classification": "classifications", } if any(context[key] not in selector[selector_keys[key]] for key in selector_keys): continue rules = {rule["field"]: rule for rule in policy["resource_rules"]["field_rules"]} if any(field not in rules or rules[field]["effect"] != "allow" for field in context["requested_fields"]): continue masks = {field: rules[field]["mask_ref"] for field in context["requested_fields"] if rules[field]["mask_ref"]} 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} return {"decision": "deny", "reason_code": "default_deny"} def create_grant(self, payload: Any, *, actor_uid: str) -> dict[str, Any]: body = _closed(payload, {"policy_uid", "approval_ref", "approval_digest", "subject_uid", "asset_uid", "provider", "expires_at", "decision"}, "trusted delivery grant") actor = self._actor(actor_uid) decision = _closed(body.get("decision"), {"decision", "reason_code", "policy_code", "policy_version", "projected_fields", "row_filter", "field_masks"}, "grant decision") if decision.get("decision") != "allow": raise PermissionError("trusted delivery grant requires allowed decision") provider = _text(body.get("provider"), "provider", 30) self.provider_registry.get(provider) expires_at = _time(body.get("expires_at"), "grant expiry") if expires_at <= self.now_factory().astimezone(UTC): raise ValueError("grant expiry must be in the future") record = { "uid": self.uid_factory(), "policy_uid": _uid(body.get("policy_uid"), "policy_uid"), "approval_ref": _text(body.get("approval_ref"), "approval reference", 200), "approval_digest": _digest_value(body.get("approval_digest"), "approval digest"), "subject_uid": _uid(body.get("subject_uid"), "subject_uid"), "asset_uid": _uid(body.get("asset_uid"), "asset_uid"), "provider": provider, "expires_at": expires_at.isoformat(), "status": "active", "reason_code": "approval_bound", "decision_digest": _digest(decision), "current_version": 1, "created_by": actor, "created_at": self.now_factory().isoformat(), } saved = self.repository.create_grant(record) result = self.repository.get_grant(saved["uid"]) if hasattr(self.repository, "get_grant") else saved self._event("grant_created", actor, {"grant_uid": result["uid"], "provider": provider, "decision_digest": record["decision_digest"]}) return result def dispatch_grant(self, grant_uid: str, *, actor_uid: str) -> dict[str, Any]: actor = self._actor(actor_uid) if hasattr(self.repository, "get_grant"): grant = self.repository.get_grant(_uid(grant_uid, "grant_uid")) elif hasattr(self.repository, "grants"): grant = copy.deepcopy(self.repository.grants.get(_uid(grant_uid, "grant_uid"))) else: raise RuntimeError("durable dispatch repository is required") if not grant or grant["status"] != "active": raise LookupError("active trusted delivery grant was not found") response = self.provider_registry.get(grant["provider"]).apply({ "schema": "dataops.trusted-delivery.v1", "grant_uid": grant["uid"], "provider": grant["provider"], "subject_uid": grant["subject_uid"], "asset_uid": grant["asset_uid"], "expires_at": grant["expires_at"], "decision_digest": grant["decision_digest"], }) response = _closed(response, {"status", "receipt_code", "response_digest"}, "provider receipt") if response.get("status") != "applied": raise RuntimeError("trusted delivery provider failed closed") 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")} self._event("grant_dispatched", actor, receipt) return receipt def provision(self, payload: Any, *, actor_uid: str, subject_roles: list[str] | None = None) -> dict[str, Any]: """Create and apply one approval-bound, digest-only capability grant.""" 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") actor = self._actor(actor_uid) provider = _text(body.get("provider"), "provider", 30) adapter = self.provider_registry.get(provider) capability = _closed(body.get("target_capability"), {"actions", "fields"}, "target capability") target_capability = { "actions": _simple_values(capability.get("actions"), "capability actions", _ACTIONS), "fields": sorted({_field(item, "capability field") for item in _values(capability.get("fields"), "capability fields", maximum=500)}), } approval_ref = _text(body.get("approval_ref"), "approval reference", 200) approval_digest = _digest_value(body.get("approval_digest"), "approval digest") policy_uid = _uid(body.get("policy_uid"), "policy_uid") subject_uid = _uid(body.get("subject_uid"), "subject_uid") asset_uid = _uid(body.get("asset_uid"), "asset_uid") context = { "subject_uid": subject_uid, "roles": list(subject_roles or []), "business_domain_uid": _uid(body.get("business_domain_uid"), "business_domain_uid"), "asset_uid": asset_uid, "purpose": _text(body.get("purpose"), "purpose", 100), "environment": _text(body.get("environment"), "environment", 30), "action": _text(body.get("action"), "action", 20), "classification": _text(body.get("classification"), "classification", 30), "requested_fields": sorted({_field(item, "requested field") for item in _values(body.get("requested_fields"), "requested fields", maximum=500)}), } # Re-evaluate from the durable, active policy version. A caller's # policy UID, capability and approval labels are only references; none # of them is allowed to enlarge the authorisation projection. decision = self.evaluate_gateway(context) active = next((item for item in self.repository.active_policies() if item.get("uid") == policy_uid), None) if active is None or decision.get("decision") != "allow" or active["code"] != decision.get("policy_code") or active["version"] != decision.get("policy_version"): raise PermissionError("trusted delivery provision is not allowed by active policy") if set(target_capability["actions"]) != {context["action"]} or set(target_capability["fields"]) != set(decision["projected_fields"]): raise PermissionError("target capability exceeds server policy projection") expires_at = _time(body.get("expires_at"), "grant expiry") if expires_at <= self.now_factory().astimezone(UTC): raise ValueError("grant expiry must be in the future") policy_expiry = _time(active["selector"]["expires_at"], "policy expiry") if expires_at > policy_expiry: raise PermissionError("grant expiry exceeds active policy") if "export" in target_capability["actions"] and self.repository.active_hold_for_asset(asset_uid): raise PermissionError("active legal hold blocks export provision") projection = { "fields": decision["projected_fields"], "row_predicate": decision["row_filter"], "masking": decision["field_masks"], } decision_digest = _digest({"policy_uid": policy_uid, "policy_code": active["code"], "policy_version": active["version"], "context": context, "projection": projection, "approval_digest": approval_digest}) record = { "uid": self.uid_factory(), "policy_uid": policy_uid, "approval_ref": approval_ref, "approval_digest": approval_digest, "subject_uid": subject_uid, "asset_uid": asset_uid, "purpose": context["purpose"], "environment": context["environment"], "target_capability": target_capability, "provider": provider, "idempotency_key": _text(body.get("idempotency_key"), "idempotency key", 160), "expires_at": expires_at.isoformat(), "status": "active", "reason_code": "approval_bound", "request_digest": _digest({"decision_digest": decision_digest, "provider": provider, "expires_at": expires_at.isoformat(), "idempotency_key": body.get("idempotency_key")}), "decision_digest": decision_digest, "current_version": 1, "created_by": actor, } if hasattr(self.repository, "enqueue_grant"): grant = self.repository.enqueue_grant(record) grant_uid = grant["uid"] if hasattr(self.repository, "get_grant"): stored = self.repository.get_grant(grant_uid) if stored and stored.get("status") == "reclaimed": raise RuntimeError("trusted delivery grant was reclaimed") if hasattr(self.repository, "get_provision_receipt"): prior_receipt = self.repository.get_provision_receipt(grant_uid, record["idempotency_key"]) if prior_receipt: return {"grant_uid": grant_uid, "status": "applied", "receipt": prior_receipt} else: raise RuntimeError("durable provision repository is required") # Persist a fenced, idempotent outbox intent before the adapter is # permitted to see the request. A process crash now leaves a retryable # pending record rather than an unprovable external side effect. if hasattr(self.repository, "prepare_provision_operation"): operation = self.repository.prepare_provision_operation( grant_uid, record["idempotency_key"], record["request_digest"] ) if operation.get("completed"): return {"grant_uid": grant_uid, "status": "applied", "receipt": operation["receipt"]} elif hasattr(self.repository, "session"): raise RuntimeError("durable provision outbox repository is required") 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"]} claim = None if hasattr(self.repository, "claim_provision_operation"): worker = f"provision:{self.uid_factory()}" claim = self.repository.claim_provision_operation(operation["delivery_uid"], worker) if claim.get("state") == "completed": return {"grant_uid": grant_uid, "status": "applied", "receipt": claim["receipt"]} if claim.get("state") != "claimed": # Another process owns the external side effect. Returning a # pending state is intentionally safer than performing a # duplicate call while that DB-time lease is alive. return {"grant_uid": grant_uid, "status": "pending", "retryable": True} try: response = _closed(adapter.apply(copy.deepcopy(envelope)), {"status", "receipt_code", "response_digest"}, "provider receipt") if response.get("status") != "applied": raise RuntimeError("trusted delivery provider failed closed") except Exception: if claim is not None: # Commit the retryable outbox state before the HTTP boundary # rolls back the failed request transaction. self.repository.fail_provision_operation(operation["delivery_uid"], worker, claim["lease_fence"]) self.repository.session.commit() raise 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")})} if claim is not None: saved = self.repository.complete_provision_operation( grant_uid, operation["delivery_uid"], worker, claim["lease_fence"], receipt ) receipt = {**receipt, "grant_uid": saved.get("uid", grant_uid)} elif hasattr(self.repository, "record_provision_receipt"): receipt = self.repository.record_provision_receipt(grant_uid, record["idempotency_key"], record["request_digest"], receipt) self._event("provision_applied", actor, {"grant_uid": grant_uid, "provider": provider, "request_digest": record["request_digest"], "response_digest": receipt["response_digest"]}) return {"grant_uid": grant_uid, "status": "applied", "receipt": receipt} def create_legal_hold(self, payload: Any, *, actor_uid: str) -> dict[str, Any]: body = _closed(payload, {"asset_uid", "operation", "approver_refs", "evidence_digest"}, "legal hold") actor = self._actor(actor_uid) asset_uid = _uid(body.get("asset_uid"), "asset_uid") if self.repository.active_hold_for_asset(asset_uid): raise RuntimeError("active legal hold already exists") operation = _text(body.get("operation"), "legal hold operation", 30) if operation not in {"freeze", "export_review", "destruction_review"}: raise ValueError("unsupported legal hold operation") approvers = sorted({_text(item, "approver reference", 200) for item in _values(body.get("approver_refs"), "approver references", minimum=2, maximum=2)}) if len(approvers) != 2: raise ValueError("legal hold requires two distinct approver references") record = { "uid": self.uid_factory(), "asset_uid": asset_uid, "operation": operation, # The SQL schema intentionally stores the two approvals in separate # columns so a parameter mapping cannot silently collapse them. "approver_ref_one": approvers[0], "approver_ref_two": approvers[1], "approver_refs": approvers, "evidence_digest": _digest_value(body.get("evidence_digest"), "evidence digest"), "status": "active", "created_by": actor, "created_at": self.now_factory().isoformat(), "current_version": 1, } saved = self.repository.create_legal_hold(record) self._event("legal_hold_active", actor, {"hold_uid": saved["uid"], "operation": operation}) return saved def release_legal_hold(self, payload: Any, *, actor_uid: str) -> dict[str, Any]: body = _closed(payload, {"hold_uid", "approval_ref", "approval_digest", "idempotency_key"}, "legal hold release") actor = self._actor(actor_uid) hold_uid = _uid(body.get("hold_uid"), "hold_uid") approval_ref = _text(body.get("approval_ref"), "approval reference", 200) approval_digest = _digest_value(body.get("approval_digest"), "approval digest") key = _text(body.get("idempotency_key"), "idempotency key", 160) request_digest = _digest({"hold_uid": hold_uid, "approval_ref": approval_ref, "approval_digest": approval_digest}) if hasattr(self.repository, "release_legal_hold"): saved = self.repository.release_legal_hold(hold_uid, key, request_digest, actor, approval_ref, approval_digest) elif hasattr(self.repository, "holds"): releases = getattr(self.repository, "hold_releases", {}) existing = releases.get(key) if existing: if existing["request_digest"] != request_digest: raise RuntimeError("trusted delivery idempotency conflict") return copy.deepcopy(existing["result"]) hold = self.repository.holds.get(hold_uid) if not hold: raise LookupError("trusted delivery legal hold was not found") if hold["created_by"] == actor or approval_ref in hold["approver_refs"]: raise PermissionError("legal hold release requires an independent approval") hold["status"] = "released" hold["current_version"] += 1 saved = {"uid": hold_uid, "status": "released", "current_version": hold["current_version"]} releases[key] = {"request_digest": request_digest, "result": copy.deepcopy(saved)} self.repository.hold_releases = releases else: raise RuntimeError("durable legal hold repository is required") self._event("legal_hold_released", actor, {"hold_uid": hold_uid, "approval_digest": approval_digest}) return saved def execute_reclaim(self, payload: Any, *, actor_uid: str) -> dict[str, Any]: """Request provider revocation first; never delete or reclaim optimistically.""" body = _closed(payload, {"grant_uid", "idempotency_key"}, "trusted delivery reclaim") actor = self._actor(actor_uid) grant_uid = _uid(body.get("grant_uid"), "grant_uid") key = _text(body.get("idempotency_key"), "idempotency key", 160) if not hasattr(self.repository, "get_grant"): raise RuntimeError("durable reclaim repository is required") grant = self.repository.get_grant(grant_uid) if not grant: raise LookupError("trusted delivery grant was not found") now = self.now_factory().astimezone(UTC) if grant.get("status") == "reclaimed": return {"grant_uid": grant_uid, "status": "reclaimed", "idempotent": True} if grant.get("status") != "active" or _time(grant["expires_at"], "grant expiry") > now: raise ValueError("trusted delivery grant is not reclaimable") if self.repository.active_hold_for_asset(grant["asset_uid"]): raise PermissionError("active legal hold blocks reclaim") expiry = _time(grant["expires_at"], "grant expiry").isoformat() request_digest = _digest({"operation": "revoke", "grant_uid": grant_uid, "provider": grant["provider"], "expires_at": expiry}) worker_id = f"reclaim-{uuid.uuid4()}" claim = None if hasattr(self.repository, "claim_reclaim_operation"): claim = self.repository.claim_reclaim_operation(grant_uid, key, request_digest, worker_id) if claim["state"] == "completed": return {"grant_uid": grant_uid, "status": "reclaimed", "idempotent": True} if claim["state"] == "leased": raise RuntimeError("trusted delivery reclaim is already leased") try: adapter = self.provider_registry.get(grant["provider"]) except PermissionError: self._event("reclaim_provider_unavailable", actor, {"grant_uid": grant_uid, "provider": grant["provider"], "request_digest": request_digest}) raise envelope = {"schema": "dataops.trusted-delivery.v2", "operation": "revoke", "grant_uid": grant_uid, "provider": grant["provider"], "asset_uid": grant["asset_uid"], "request_digest": request_digest} revoke = getattr(adapter, "revoke", None) if not callable(revoke): self._event("reclaim_provider_unavailable", actor, {"grant_uid": grant_uid, "provider": grant["provider"], "request_digest": request_digest}) raise PermissionError("trusted delivery provider revoke is disabled") try: response = _closed(revoke(copy.deepcopy(envelope)), {"status", "receipt_code", "response_digest"}, "provider receipt") except Exception: if claim: self.repository.fail_reclaim_operation(claim["delivery_uid"], worker_id, claim["lease_fence"]) raise if response.get("status") != "revoked": if claim: self.repository.fail_reclaim_operation(claim["delivery_uid"], worker_id, claim["lease_fence"]) raise RuntimeError("trusted delivery provider failed closed") 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")})} if claim: saved = self.repository.complete_claimed_reclaim(grant_uid, claim["delivery_uid"], worker_id, claim["lease_fence"], receipt) elif hasattr(self.repository, "complete_reclaim"): saved = self.repository.complete_reclaim(grant_uid, key, request_digest, receipt) else: saved = self.repository.update_grant_status(grant_uid, grant["current_version"], "reclaimed", "expired_reclaimed") self._event("grant_reclaimed", actor, {"grant_uid": grant_uid, "provider": grant["provider"], "request_digest": request_digest, "response_digest": receipt["response_digest"]}) return {"grant_uid": grant_uid, "status": saved["status"], "receipt": receipt} def preview_reclaim(self, *, as_of: Any, actor_uid: str) -> list[dict[str, Any]]: """Return eligible identifiers only; this query has no state transition.""" self._actor(actor_uid) instant = _time(as_of, "as_of") candidates = [] for grant in self.repository.grants_due_for_revoke(instant): if grant.get("status") != "active" or _time(grant["expires_at"], "grant expiry") > instant: continue if self.repository.active_hold_for_asset(grant["asset_uid"]): continue candidates.append({"grant_uid": grant["uid"], "provider": grant.get("provider", "unknown"), "status": "eligible"}) return candidates def revoke_expired(self, *, as_of: datetime, actor_uid: str) -> list[dict[str, Any]]: actor = self._actor(actor_uid) instant = _time(as_of, "as_of") revoked = [] for grant in self.repository.grants_due_for_revoke(instant): if _time(grant["expires_at"], "grant expiry") > instant or self.repository.active_hold_for_asset(grant["asset_uid"]): continue saved = self.repository.update_grant_status(grant["uid"], grant["current_version"], "revoked", "expired") self._event("grant_revoked", actor, {"grant_uid": saved["uid"], "reason_code": "expired"}) revoked.append(saved) return revoked @staticmethod def mask_preview(payload: Any) -> dict[str, Any]: body = _closed(payload, {"mode", "mask_ref", "value"}, "mask preview") mode = _text(body.get("mode"), "mask mode", 20) if mode not in {"static", "dynamic", "display", "export"}: raise ValueError("unsupported mask mode") mask_ref = _text(body.get("mask_ref"), "mask ref", 120) value = _text(body.get("value"), "mask value", 1000) return {"mode": mode, "mask_ref": mask_ref, "masked_value": "***-****", "token_digest": _digest({"mask_ref": mask_ref, "value": value})}