"""Fail-closed enterprise identity policy, mapping, lifecycle and emergency access.""" from __future__ import annotations import hashlib import ipaddress import json import os import re from collections.abc import Callable, Mapping from dataclasses import dataclass, field from datetime import UTC, datetime, timedelta from typing import Any from urllib.parse import urlsplit from app.core.common.identifiers import new_governance_uid class IdentityPolicyError(ValueError): pass class IdentityUpstreamError(IdentityPolicyError): """Safe boundary error for unavailable or malformed enterprise IdP responses.""" pass _SENSITIVE_KEYS = {"code", "token", "id_token", "access_token", "refresh_token", "verifier", "code_verifier", "secret", "client_secret", "claims", "cookie", "password", "assertion", "authorization_url"} def _utc(value: datetime) -> datetime: if value.tzinfo is None or value.utcoffset() is None: raise IdentityPolicyError("timestamp must be timezone-aware") return value.astimezone(UTC) def _safe_https_url(value: str, *, redirect: bool = False) -> None: parsed = urlsplit(value) if parsed.scheme != "https" or not parsed.hostname or parsed.username or parsed.password: raise IdentityPolicyError("identity endpoints require absolute HTTPS URLs") if parsed.fragment or (redirect and parsed.query): raise IdentityPolicyError("redirect and identity URLs may not contain query/fragment") host = parsed.hostname.lower() if host == "localhost" or host.endswith(".localhost"): raise IdentityPolicyError("localhost identity endpoints are forbidden") try: address = ipaddress.ip_address(host) except ValueError: return if not address.is_global: raise IdentityPolicyError("private or reserved identity endpoint is forbidden") @dataclass(frozen=True) class IdpConfig: provider_uid: str version: int issuer: str client_id: str secret_ref: str = field(repr=False) authorization_endpoint: str token_endpoint: str jwks_uri: str redirect_uris: tuple[str, ...] post_login_redirect_uris: tuple[str, ...] algorithms: tuple[str, ...] mapping_version: str status: str = "draft" def validate(self) -> IdpConfig: if self.version < 1 or not self.provider_uid or not self.client_id or not self.mapping_version: raise IdentityPolicyError("incomplete IdP configuration") for endpoint in (self.issuer, self.authorization_endpoint, self.token_endpoint, self.jwks_uri): _safe_https_url(endpoint) if not self.redirect_uris: raise IdentityPolicyError("redirect allowlist is required") for uri in self.redirect_uris: _safe_https_url(uri, redirect=True) if not self.post_login_redirect_uris: raise IdentityPolicyError("post-login redirect allowlist is required") for uri in self.post_login_redirect_uris: _safe_https_url(uri, redirect=True) if not re.fullmatch(r"env:DATAOPS_OIDC_[A-Z0-9_]+", self.secret_ref): raise IdentityPolicyError("client secret must use a dedicated DATAOPS_OIDC environment reference") if not self.algorithms or any(alg not in {"RS256", "RS384", "RS512", "ES256", "ES384"} for alg in self.algorithms): raise IdentityPolicyError("only explicitly configured asymmetric algorithms are allowed") return self def resolve_secret(self) -> str: self.validate() value = os.environ.get(self.secret_ref[4:]) if not value: raise IdentityPolicyError("configured secret reference is unavailable") return value def public_dict(self) -> dict[str, Any]: return {"provider_uid": self.provider_uid, "version": self.version, "issuer": self.issuer, "client_id": self.client_id, "authorization_endpoint": self.authorization_endpoint, "redirect_uris": list(self.redirect_uris), "algorithms": list(self.algorithms), "post_login_redirect_uris": list(self.post_login_redirect_uris), "mapping_version": self.mapping_version, "status": self.status, "secret_ref": "env:DATAOPS_OIDC_***"} def sanitize_audit_detail(value: Any) -> Any: if isinstance(value, Mapping): safe: dict[str, Any] = {} for key, item in value.items(): lowered = str(key).lower() if lowered in _SENSITIVE_KEYS or any(part in lowered for part in ("secret", "token", "verifier")): continue if lowered.endswith("claims"): continue safe[str(key)] = sanitize_audit_detail(item) return safe if isinstance(value, (list, tuple)): return [sanitize_audit_detail(item) for item in value] return value @dataclass(frozen=True) class MappedIdentity: subject: str username: str display_name: str department: str groups: tuple[str, ...] roles: tuple[str, ...] business_domain_uids: tuple[str, ...] object_types: tuple[str, ...] environments: tuple[str, ...] data_scopes: tuple[str, ...] mapping_version: str claims_digest: str evidence: Mapping[str, Any] = field(repr=False) class ClaimMapper: REQUIRED = ("sub", "preferred_username", "name", "department", "groups") RULE_FIELDS = {"roles", "business_domain_uids", "object_types", "environments", "data_scopes"} def __init__(self, *, version: str, group_rules: Mapping[str, Mapping[str, Any]]) -> None: if not isinstance(version, str) or not version.strip(): raise IdentityPolicyError("mapping version is required") if not isinstance(group_rules, Mapping) or not group_rules: raise IdentityPolicyError("at least one claims group rule is required") normalized: dict[str, dict[str, tuple[str, ...]]] = {} for group, raw_rule in group_rules.items(): if not isinstance(group, str) or not group.strip() or not isinstance(raw_rule, Mapping): raise IdentityPolicyError("claims group rules must be named objects") if set(raw_rule) - self.RULE_FIELDS: raise IdentityPolicyError("claims group rule contains unsupported fields") rule: dict[str, tuple[str, ...]] = {} for field_name in self.RULE_FIELDS: values = raw_rule.get(field_name, ()) if not isinstance(values, (list, tuple)) or isinstance(values, (str, bytes)): raise IdentityPolicyError("claims mapping scope fields must be string arrays") if any(not isinstance(value, str) or not value.strip() for value in values): raise IdentityPolicyError("claims mapping scope values must be non-empty strings") rule[field_name] = tuple(sorted(set(values))) if not rule["roles"] or not set(rule["roles"]).issubset({"admin", "editor", "viewer"}): raise IdentityPolicyError("mapping contains an unsupported or empty platform role") normalized[group.strip()] = rule self.version = version.strip() self.group_rules = normalized def map(self, claims: Mapping[str, Any]) -> MappedIdentity: if not isinstance(claims, Mapping): raise IdentityPolicyError("identity claims must be an object") limits = {"sub": 500, "preferred_username": 64, "name": 300, "department": 300} if any(not isinstance(claims.get(key), str) or not claims[key].strip() or len(claims[key].strip()) > maximum for key, maximum in limits.items()): raise IdentityPolicyError("required identity claim is missing") raw_groups = claims.get("groups") if (not isinstance(raw_groups, (list, tuple)) or isinstance(raw_groups, (str, bytes)) or not raw_groups or any(not isinstance(item, str) or not item.strip() or len(item.strip()) > 300 for item in raw_groups)): raise IdentityPolicyError("identity groups claim must be a non-empty string array") groups = tuple(sorted({item.strip() for item in raw_groups})) matched = [self.group_rules[group] for group in groups if group in self.group_rules] if not matched: raise IdentityPolicyError("claims do not match an approved mapping rule") def collect(key: str) -> tuple[str, ...]: return tuple(sorted({str(value) for rule in matched for value in rule.get(key, ())})) roles = collect("roles") if not roles or not set(roles).issubset({"admin", "editor", "viewer"}): raise IdentityPolicyError("mapping resolves to no approved platform role") canonical = json.dumps(claims, ensure_ascii=False, sort_keys=True, separators=(",", ":"), default=str) digest = hashlib.sha256(canonical.encode()).hexdigest() evidence = {"mapping_version": self.version, "claims_digest": digest, "matched_group_digests": [hashlib.sha256(group.encode()).hexdigest() for group in groups if group in self.group_rules]} return MappedIdentity(claims["sub"].strip(), claims["preferred_username"].strip(), claims["name"].strip(), claims["department"].strip(), groups, roles, collect("business_domain_uids"), collect("object_types"), collect("environments"), collect("data_scopes"), self.version, digest, evidence) class DirectorySynchronizer: EVENTS = {"JOINER", "MOVER", "LEAVER", "DISABLE", "RESTORE"} ORGANIZATION_EVENTS = {"DEPARTMENT", "GROUP"} ORGANIZATION_ACTIONS = {"UPSERT", "DISABLE", "RESTORE"} def __init__(self, repository: Any, sessions: Any, *, clock: Callable[[], datetime] | None = None) -> None: self.repository = repository self.sessions = sessions self.clock = clock or (lambda: datetime.now(UTC)) def apply(self, *, provider_uid: str, source: str, source_event_id: str, cursor: str, cursor_sequence: int, event_type: str, attributes: Mapping[str, Any], subject: str | None = None) -> dict[str, Any]: event_type = event_type.upper() if event_type not in self.EVENTS | self.ORGANIZATION_EVENTS or not all((provider_uid, source, source_event_id, cursor)): raise IdentityPolicyError("invalid directory delta event") if event_type in self.EVENTS and not subject: raise IdentityPolicyError("person lifecycle event requires enterprise subject") if not isinstance(cursor_sequence, int) or cursor_sequence < 1: raise IdentityPolicyError("directory cursor sequence must be a positive integer") digest = hashlib.sha256(json.dumps({"cursor": cursor, "event_type": event_type, "subject": subject, "attributes": attributes}, sort_keys=True, default=str).encode()).hexdigest() event = {"provider_uid": provider_uid, "source": source, "source_event_id": source_event_id, "cursor": cursor, "cursor_sequence": cursor_sequence, "event_type": event_type, "subject": subject, "payload_digest": digest, "processed_at": _utc(self.clock()).isoformat()} def mutate_organization() -> dict[str, Any]: node_type = event_type.lower() external_id = str(attributes.get("external_id") or "").strip() action = str(attributes.get("lifecycle_action") or attributes.get("action") or "UPSERT").upper() if not external_id or action not in self.ORGANIZATION_ACTIONS: raise IdentityPolicyError("invalid organization node delta") existing = self.repository.get_organization_node(provider_uid, node_type, external_id, for_update=True) if action in {"DISABLE", "RESTORE"} and not existing: raise IdentityPolicyError("organization node lifecycle target does not exist") display_name = str(attributes.get("display_name") or (existing or {}).get("display_name") or "").strip() if not display_name: raise IdentityPolicyError("organization node display name is required") node = { "uid": (existing or {}).get("uid") or new_governance_uid(), "provider_uid": provider_uid, "external_id": external_id, "node_type": node_type, "parent_external_id": attributes.get("parent_external_id"), "display_name": display_name, "status": "disabled" if action == "DISABLE" else "active", "attributes": dict(attributes.get("attributes", {})), "updated_at": _utc(self.clock()), } self.repository.put_organization_node(node, commit=False) return {"node_uid": node["uid"], "node_type": node_type, "external_id": external_id, "status": node["status"], "lifecycle_action": action, "idempotent": False} def mutate_identity() -> dict[str, Any]: identity = self.repository.get_identity(provider_uid, str(subject), for_update=True) if not identity and event_type != "JOINER": raise IdentityPolicyError("person lifecycle target does not exist") identity = identity or {"user_uid": new_governance_uid(), "token_version": 0} old_roles = tuple(identity.get("roles", ())) roles = sorted(set(attributes.get("roles", ()))) if not set(roles).issubset({"admin", "editor", "viewer"}): raise IdentityPolicyError("directory event contains an unsupported platform role") username = attributes.get("username") or identity.get("username") if not username: raise IdentityPolicyError("person lifecycle username is required") identity.update(subject=subject, username=username, department=attributes.get("department", identity.get("department", "")), display_name=attributes.get("display_name") or identity.get("display_name") or username, provider_uid=provider_uid, groups=sorted(attributes.get("groups", identity.get("groups", ()))), roles=roles, mapping_version=attributes.get("mapping_version") or identity.get("mapping_version"), claims_digest=attributes.get("claims_digest") or identity.get("claims_digest"), authorization_scope=dict(attributes.get("authorization_scope", identity.get("authorization_scope", {}))), status="disabled" if event_type in {"LEAVER", "DISABLE"} else "active", token_version=int(identity.get("token_version", 0)) + 1, updated_at=_utc(self.clock()).isoformat()) # RESTORE deliberately recalculates from current input and never restores old roles. self.repository.put_identity(provider_uid, str(subject), identity, commit=False) if event_type in {"MOVER", "LEAVER", "DISABLE", "RESTORE"} or old_roles != tuple(identity["roles"]): self.repository.revoke_subject_sessions(provider_uid, str(subject), reason=f"directory_{event_type.lower()}", commit=False) return {"user_uid": identity["user_uid"], "status": identity["status"], "token_version": identity["token_version"], "idempotent": False} mutation = mutate_organization if event_type in self.ORGANIZATION_EVENTS else mutate_identity return self.repository.apply_directory_event_atomic(event, mutation) class EmergencyAccess: MAX_DURATION = timedelta(hours=2) def __init__(self, repository: Any, *, clock: Callable[[], datetime] | None = None) -> None: self.repository = repository self.clock = clock or (lambda: datetime.now(UTC)) def request(self, *, requester_uid: str, account_uid: str, reason: str, expires_at: datetime, account_is_local_active_admin: bool) -> dict[str, Any]: now = _utc(self.clock()) expiry = _utc(expires_at) if (requester_uid == account_uid or not reason.strip() or not now < expiry or expiry - now > self.MAX_DURATION or not account_is_local_active_admin): raise IdentityPolicyError("invalid emergency access request") record = {"uid": new_governance_uid(), "requester_uid": requester_uid, "account_uid": account_uid, "reason": reason.strip(), "status": "pending", "approver_uids": [], "requested_at": now, "expires_at": expiry, "reviewed_at": None} self.repository.put_emergency(record) return record def approve(self, request_uid: str, *, approver_uid: str) -> dict[str, Any]: record = self.repository.get_emergency(request_uid) if not record or record["status"] != "pending" or _utc(self.clock()) >= record["expires_at"]: raise IdentityPolicyError("emergency request is not approvable") forbidden = {record["requester_uid"], record["account_uid"], *record["approver_uids"]} if approver_uid in forbidden: raise IdentityPolicyError("emergency access requires two distinct non-self approvers") record["approver_uids"].append(approver_uid) self.repository.put_emergency(record) return record def activate(self, request_uid: str) -> dict[str, Any]: record = self.repository.get_emergency(request_uid) if not record or record["status"] != "pending" or len(set(record["approver_uids"])) != 2 or _utc(self.clock()) >= record["expires_at"]: raise IdentityPolicyError("two approvals in the active time window are required") record.update(status="active", activated_at=_utc(self.clock()), alert_required=True) self.repository.put_emergency(record) return record def review(self, request_uid: str, *, reviewer_uid: str, outcome: str) -> dict[str, Any]: record = self.repository.get_emergency(request_uid) if record and record["status"] == "active" and _utc(self.clock()) >= record["expires_at"]: record["status"] = "expired" self.repository.put_emergency(record) if not record or record["status"] not in {"closed", "expired"} or reviewer_uid == record["account_uid"]: raise IdentityPolicyError("invalid emergency review") record.update(status="reviewed", reviewed_at=_utc(self.clock()), reviewer_uid=reviewer_uid, review_outcome=outcome) self.repository.put_emergency(record) return record def close(self, request_uid: str, *, actor_uid: str) -> dict[str, Any]: record = self.repository.get_emergency(request_uid) if record and record["status"] == "active" and _utc(self.clock()) >= record["expires_at"]: record["status"] = "expired" self.repository.put_emergency(record) raise IdentityPolicyError("emergency access already expired") if not record or record["status"] != "active" or actor_uid == record["account_uid"]: raise IdentityPolicyError("invalid emergency access closure") record.update(status="closed", closed_at=_utc(self.clock()), closed_by_uid=actor_uid) self.repository.put_emergency(record) return record