| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345 |
- """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
|