"""Enterprise identity APIs. Public routes are explicitly enumerated in policy.""" from __future__ import annotations import hashlib import json import secrets import uuid from datetime import UTC, datetime, timedelta from functools import wraps from urllib.parse import urlsplit from flask import current_app, g, jsonify, redirect, request from itsdangerous import BadSignature, URLSafeTimedSerializer from sqlalchemy import text from sqlalchemy.exc import IntegrityError, StatementError from app import db from app.api.system import bp from app.core.common.identifiers import new_governance_uid from app.core.system.enterprise_identity import ( ClaimMapper, DirectorySynchronizer, EmergencyAccess, IdentityPolicyError, IdentityUpstreamError, IdpConfig, sanitize_audit_detail, ) from app.core.system.identity_repository import PostgresIdentityRepository from app.core.system.identity_sessions import SessionManager from app.core.system.oidc import OidcClient, OidcTransport from app.core.system.permissions import ( IDENTITY_MANAGE, IDENTITY_OPERATE, IDENTITY_READ, require_permissions, ) from app.core.system.tokens import issue_access_token from app.models.result import failed, success def _body() -> dict: value = request.get_json(silent=True) if not isinstance(value, dict): raise IdentityPolicyError("JSON object body is required") return value def _required_str(data: dict, key: str, *, maximum: int = 500) -> str: value = data.get(key) if not isinstance(value, str) or not value.strip() or len(value.strip()) > maximum: raise IdentityPolicyError(f"{key} is required or invalid") return value.strip() def _optional_str(data: dict, key: str, default: str, *, maximum: int) -> str: if key not in data: return default return _required_str(data, key, maximum=maximum) def _uuid(value, field: str) -> str: try: return str(uuid.UUID(str(value))) except (ValueError, TypeError, AttributeError) as exc: raise IdentityPolicyError(f"{field} must be a UUID") from exc def _positive_int(value, field: str, *, maximum: int = 2_147_483_647) -> int: if isinstance(value, bool): raise IdentityPolicyError(f"{field} must be a positive integer") try: result = int(value) except (ValueError, TypeError) as exc: raise IdentityPolicyError(f"{field} must be a positive integer") from exc if result < 1 or result > maximum: raise IdentityPolicyError(f"{field} must be a positive integer") return result def _timestamp(value, field: str) -> datetime: if not isinstance(value, str): raise IdentityPolicyError(f"{field} must be an ISO-8601 timestamp") try: parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) except ValueError as exc: raise IdentityPolicyError(f"{field} must be an ISO-8601 timestamp") from exc if parsed.tzinfo is None or parsed.utcoffset() is None: raise IdentityPolicyError(f"{field} must include a timezone") return parsed def _string_array(data: dict, key: str, *, default=None) -> tuple[str, ...]: value = data.get(key, default) if not isinstance(value, list) or not value or any(not isinstance(item, str) or not item.strip() for item in value): raise IdentityPolicyError(f"{key} must be a non-empty string array") return tuple(item.strip() for item in value) def _validate_mapping(version: str, group_rules) -> ClaimMapper: if not isinstance(group_rules, dict): raise IdentityPolicyError("group_rules must be an object") return ClaimMapper(version=version, group_rules=group_rules) def _identity_input_guard(view): @wraps(view) def guarded(*args, **kwargs): try: return view(*args, **kwargs) except IdentityPolicyError: raise except (KeyError, TypeError, ValueError, IntegrityError, StatementError) as exc: db.session.rollback() raise IdentityPolicyError("invalid enterprise identity request") from exc return guarded def _repo() -> PostgresIdentityRepository: return PostgresIdentityRepository(db.session) def _config(row) -> IdpConfig: return IdpConfig(provider_uid=str(row["provider_uid"]), version=row["version"], issuer=row["issuer"], client_id=row["client_id"], secret_ref=row["secret_ref"], authorization_endpoint=row["authorization_endpoint"], token_endpoint=row["token_endpoint"], jwks_uri=row["jwks_uri"], redirect_uris=tuple(row["redirect_uris"]), post_login_redirect_uris=tuple(dict(row["config"]).get("post_login_redirect_uris", ())), algorithms=tuple(row["algorithms"]), mapping_version=row["mapping_version"], status=row["status"]) def _load_config(provider_uid: str, *, version: int | None = None, active: bool = True): clause = "AND status='active'" if active else "" version_clause = "AND version=:version" if version is not None else "" return db.session.execute(text(f""" SELECT provider_uid::text,version,status,issuer,client_id,secret_ref,authorization_endpoint, token_endpoint,jwks_uri,redirect_uris,algorithms,mapping_version,config FROM public.identity_idp_config_versions WHERE provider_uid=CAST(:provider AS uuid) {version_clause} {clause} ORDER BY version DESC LIMIT 1 """), {"provider": provider_uid, "version": version}).mappings().one_or_none() def _audit(event_type: str, outcome: str, *, actor_uid: str | None = None, provider_uid: str | None = None, resource_uid: str | None = None, enterprise_subject: str | None = None, detail: dict | None = None) -> None: _repo().audit({"event_type": event_type, "outcome": outcome, "actor_uid": actor_uid, "provider_uid": provider_uid, "resource_type": "enterprise_identity", "resource_uid": resource_uid, "subject_digest": (hashlib.sha256(enterprise_subject.encode()).hexdigest() if enterprise_subject else None), "safe_detail": sanitize_audit_detail(detail or {})}) def _cookie_secure() -> bool: return bool(current_app.config.get("IDENTITY_COOKIE_SECURE", True)) def _origin(uri: str) -> str: parts = urlsplit(uri) return f"{parts.scheme}://{parts.netloc}" def _browser_origin_allowed(origin: str) -> bool: rows = db.session.execute(text("SELECT config FROM public.identity_idp_config_versions WHERE status='active'")).scalars().all() allowed = {_origin(uri) for config in rows for uri in dict(config).get("post_login_redirect_uris", ())} allowed.update(current_app.config.get("IDENTITY_ALLOWED_TEST_ORIGINS", ())) return origin in allowed @bp.errorhandler(IdentityPolicyError) def identity_policy_error(exc): try: db.session.rollback() _audit("identity_policy_denied", "denied", actor_uid=getattr(g, "current_user", {}).get("id"), detail={"path": request.path, "reason": str(exc)[:160]}) except Exception: # pragma: no cover - failure reporting must not expose internals db.session.rollback() return jsonify(failed(str(exc), code=400)), 400 @bp.errorhandler(IdentityUpstreamError) def identity_upstream_error(exc): db.session.rollback() try: _audit("identity_upstream_failure", "failure", actor_uid=getattr(g, "current_user", {}).get("id"), detail={"path": request.path}) except Exception: # pragma: no cover db.session.rollback() return jsonify(failed("enterprise identity provider is unavailable", code=502)), 502 @bp.route("/identity/providers", methods=["GET"]) def identity_providers(): rows = db.session.execute(text(""" SELECT provider_uid::text,version,issuer,config->>'display_name' display_name FROM public.identity_idp_config_versions WHERE status='active' ORDER BY provider_uid """)).mappings().all() return jsonify(success([{"provider_uid": row["provider_uid"], "version": row["version"], "issuer": row["issuer"], "display_name": row["display_name"] or "企业单点登录"} for row in rows])) @bp.route("/identity/authorize", methods=["POST"]) @_identity_input_guard def identity_authorize(): data = _body() provider_uid = _uuid(data.get("provider_uid"), "provider_uid") row = _load_config(provider_uid) if not row: raise IdentityPolicyError("active identity provider not found") config = _config(row) callback_uri = _required_str(data, "redirect_uri", maximum=2048) post_login_uri = _required_str(data, "post_login_uri", maximum=2048) if post_login_uri not in config.post_login_redirect_uris: raise IdentityPolicyError("post-login redirect URI is not allowlisted") started = OidcClient(_repo()).begin(config, callback_uri) serializer = URLSafeTimedSerializer(current_app.config["SECRET_KEY"], salt="oidc-flow") response = jsonify(success({"authorization_url": started.url, "provider_uid": config.provider_uid})) response.set_cookie("dataops_oidc_flow", serializer.dumps({"state": started.state, "verifier": started.code_verifier, "provider_uid": config.provider_uid, "redirect_uri": callback_uri, "post_login_uri": post_login_uri}), httponly=True, secure=_cookie_secure(), samesite="Lax", max_age=300, path="/api/system/identity/callback") _audit("oidc_authorize", "success", provider_uid=config.provider_uid, detail={"redirect_uri": data["redirect_uri"]}) return response @bp.route("/identity/callback", methods=["GET"]) @_identity_input_guard def identity_callback(): serializer = URLSafeTimedSerializer(current_app.config["SECRET_KEY"], salt="oidc-flow") try: flow = serializer.loads(request.cookies.get("dataops_oidc_flow", ""), max_age=300) except BadSignature as exc: raise IdentityPolicyError("OIDC browser flow binding is invalid or expired") from exc if not isinstance(flow, dict) or request.args.get("state") != flow.get("state"): raise IdentityPolicyError("OIDC state mismatch") provider_uid = _uuid(flow.get("provider_uid"), "provider_uid") row = _load_config(provider_uid) if not row: raise IdentityPolicyError("active identity provider not found") row = dict(row) config = _config(row) config.validate() callback_uri = _required_str(flow, "redirect_uri", maximum=2048) verifier = _required_str(flow, "verifier", maximum=512) state = _required_str(flow, "state", maximum=512) post_login_uri = _required_str(flow, "post_login_uri", maximum=2048) code = request.args.get("code") if not isinstance(code, str) or not code or len(code) > 4096: raise IdentityPolicyError("OIDC authorization code is required") # Do not keep a database transaction open while calling the external IdP. db.session.rollback() transport = OidcTransport() id_token = transport.exchange_code(config, code=code, redirect_uri=callback_uri, verifier=verifier) jwks = transport.fetch_jwks(config) repo = _repo() try: claims = OidcClient(repo).verify_callback(config, state=state, redirect_uri=callback_uri, code_verifier=verifier, id_token=id_token, jwks=jwks, commit=False) mapped = ClaimMapper(version=config.mapping_version, group_rules=dict(row["config"]).get("group_rules", {})).map(claims) link = repo.get_identity(config.provider_uid, mapped.subject, for_update=True) if link and link["status"] != "active": raise IdentityPolicyError("enterprise identity is disabled") scope = {"business_domain_uids": mapped.business_domain_uids, "object_types": mapped.object_types, "environments": mapped.environments, "data_scopes": mapped.data_scopes} changed = bool(link and (sorted(link["roles"]) != sorted(mapped.roles) or dict(link["authorization_scope"]) != json.loads(json.dumps(scope)))) user_uid = link["user_uid"] if link else new_governance_uid() token_version = (int(link["token_version"]) + 1 if changed else int(link["token_version"])) if link else 1 repo.put_identity(config.provider_uid, mapped.subject, { "provider_uid": config.provider_uid, "user_uid": user_uid, "username": mapped.username, "display_name": mapped.display_name, "department": mapped.department, "groups": list(mapped.groups), "roles": list(mapped.roles), "authorization_scope": scope, "mapping_version": mapped.mapping_version, "claims_digest": mapped.claims_digest, "status": "active", "token_version": token_version, }, commit=False) if changed: repo.revoke_subject_sessions(config.provider_uid, mapped.subject, reason="claims_mapping_changed", commit=False) credentials = SessionManager(repo, secret=current_app.config["SECRET_KEY"]).create( subject=mapped.subject, user_uid=user_uid, roles=list(mapped.roles), identity_source="oidc", token_version=token_version, provider_uid=config.provider_uid, commit=False) exchange = secrets.token_urlsafe(40) repo.insert_exchange_code({"uid": new_governance_uid(), "code_hash": hashlib.sha256(exchange.encode()).hexdigest(), "session_uid": credentials.session_uid, "redirect_uri": post_login_uri, "expires_at": datetime.now(UTC) + timedelta(minutes=2)}, commit=False) db.session.commit() except IdentityPolicyError: db.session.rollback() raise except Exception as exc: db.session.rollback() raise IdentityPolicyError("OIDC callback persistence failed") from exc _audit("oidc_login", "success", provider_uid=config.provider_uid, resource_uid=user_uid, detail={"mapping_version": mapped.mapping_version, "claims_digest": mapped.claims_digest}) response = redirect(f"{flow['post_login_uri']}?exchange_code={exchange}", code=302) response.set_cookie("dataops_refresh", credentials.refresh_token, httponly=True, secure=_cookie_secure(), samesite="Strict", max_age=8 * 3600, path="/api/system/identity") response.delete_cookie("dataops_oidc_flow", path="/api/system/identity/callback") return response @bp.route("/identity/exchange", methods=["POST"]) @_identity_input_guard def identity_exchange(): data = _body() exchange_code = _required_str(data, "exchange_code", maximum=1024) redirect_uri = _required_str(data, "redirect_uri", maximum=2048) digest = hashlib.sha256(exchange_code.encode()).hexdigest() row = db.session.execute(text(""" UPDATE public.identity_exchange_codes SET status='consumed',consumed_at=CURRENT_TIMESTAMP WHERE code_hash=:digest AND status='pending' AND expires_at>CURRENT_TIMESTAMP AND redirect_uri=:redirect RETURNING session_uid::text,redirect_uri """), {"digest": digest, "redirect": redirect_uri}).mappings().one_or_none() if not row or request.headers.get("Origin") != _origin(row["redirect_uri"]): db.session.rollback() raise IdentityPolicyError("exchange code or browser origin is invalid") db.session.commit() session = _repo().get_session(row["session_uid"]) if not session or session["status"] != "active": raise IdentityPolicyError("identity session is unavailable") token = issue_access_token(user_id=session["user_uid"], roles=list(session["roles"]), secret=current_app.config["SECRET_KEY"], session_uid=session["uid"], token_version=session["token_version"], identity_source=session["identity_source"]) return jsonify(success({"token": token, "identity_source": session["identity_source"]})) @bp.route("/identity/refresh", methods=["POST"]) def identity_refresh(): origin = request.headers.get("Origin", "") if not _browser_origin_allowed(origin): raise IdentityPolicyError("browser origin is not allowlisted") credentials = SessionManager(_repo(), secret=current_app.config["SECRET_KEY"]).refresh(request.cookies.get("dataops_refresh", "")) session = _repo().get_session(credentials.session_uid) _audit("identity_refresh", "success", actor_uid=session["user_uid"] if session else None, provider_uid=session.get("provider_uid") if session else None, resource_uid=credentials.session_uid, detail={"rotated_from_uid": session.get("rotated_from_uid") if session else None}) response = jsonify(success({"token": credentials.access_token})) response.set_cookie("dataops_refresh", credentials.refresh_token, httponly=True, secure=_cookie_secure(), samesite="Strict", max_age=8 * 3600, path="/api/system/identity") return response @bp.route("/identity/logout", methods=["POST"]) @require_permissions(IDENTITY_READ) def identity_logout(): session_uid = g.current_user.get("sid") session = _repo().get_session(session_uid) if session_uid else None if session_uid: SessionManager(_repo(), secret=current_app.config["SECRET_KEY"]).revoke(session_uid) _audit("identity_logout", "success", actor_uid=g.current_user.get("id"), provider_uid=session.get("provider_uid") if session else None, resource_uid=session_uid) response = jsonify(success({"revoked": True})) response.delete_cookie("dataops_refresh", path="/api/system/identity") return response @bp.route("/identity/idp-versions", methods=["GET"]) @require_permissions(IDENTITY_MANAGE) def identity_idp_versions(): rows = db.session.execute(text("SELECT provider_uid::text,version,status,issuer,client_id,secret_ref,redirect_uris,algorithms,mapping_version,created_at FROM public.identity_idp_config_versions ORDER BY provider_uid,version DESC")).mappings().all() return jsonify(success([{**dict(row), "secret_ref": "env:DATAOPS_OIDC_***"} for row in rows])) @bp.route("/identity/idp-versions", methods=["POST"]) @require_permissions(IDENTITY_MANAGE) @_identity_input_guard def identity_idp_version_create(): data = _body() provider_uid = _uuid(data["provider_uid"], "provider_uid") if data.get("provider_uid") else new_governance_uid() mapping_version = _required_str(data, "mapping_version", maximum=120) group_rules = data.get("group_rules") _validate_mapping(mapping_version, group_rules) config = IdpConfig(provider_uid=provider_uid, version=_positive_int(data.get("version"), "version"), issuer=_required_str(data, "issuer", maximum=2048), client_id=_required_str(data, "client_id", maximum=300), secret_ref=_required_str(data, "secret_ref", maximum=160), authorization_endpoint=_required_str(data, "authorization_endpoint", maximum=2048), token_endpoint=_required_str(data, "token_endpoint", maximum=2048), jwks_uri=_required_str(data, "jwks_uri", maximum=2048), redirect_uris=_string_array(data, "redirect_uris"), post_login_redirect_uris=_string_array(data, "post_login_redirect_uris"), algorithms=_string_array(data, "algorithms", default=["RS256"]), mapping_version=mapping_version).validate() db.session.execute(text(""" INSERT INTO public.identity_idp_config_versions(uid,provider_uid,version,status,issuer,client_id,secret_ref,authorization_endpoint, token_endpoint,jwks_uri,redirect_uris,algorithms,mapping_version,config,created_by) VALUES(CAST(:uid AS uuid),CAST(:provider AS uuid),:version,'draft',:issuer,:client,:secret,:auth,:token,:jwks, CAST(:redirects AS jsonb),CAST(:algorithms AS jsonb),:mapping,CAST(:config AS jsonb),CAST(:actor AS uuid)) """), {"uid": new_governance_uid(), "provider": config.provider_uid, "version": config.version, "issuer": config.issuer, "client": config.client_id, "secret": config.secret_ref, "auth": config.authorization_endpoint, "token": config.token_endpoint, "jwks": config.jwks_uri, "redirects": json.dumps(config.redirect_uris), "algorithms": json.dumps(config.algorithms), "mapping": config.mapping_version, "config": json.dumps({"display_name": data.get("display_name"), "group_rules": group_rules, "post_login_redirect_uris": data.get("post_login_redirect_uris", [])}), "actor": g.current_user["id"]}) db.session.commit() _audit("idp_version_created", "success", actor_uid=g.current_user["id"], provider_uid=config.provider_uid) return jsonify(success(config.public_dict())), 201 @bp.route("/identity/idp-versions///", methods=["POST"]) @require_permissions(IDENTITY_MANAGE) @_identity_input_guard def identity_idp_action(provider_uid: str, version: int, action: str): provider_uid = _uuid(provider_uid, "provider_uid") version = _positive_int(version, "version") row = _load_config(provider_uid, version=version, active=False) if not row: raise IdentityPolicyError("IdP configuration version not found") config = _config(row) config.validate() _validate_mapping(config.mapping_version, dict(row["config"]).get("group_rules")) if action == "dry-run": result = {"valid": True, "secret_available": bool(config.resolve_secret())} _audit("idp_dry_run", "success", actor_uid=g.current_user["id"], provider_uid=provider_uid, resource_uid=f"{provider_uid}:{version}", detail={"version": version}) return jsonify(success(result)) if action in {"activate", "rotate"}: config.resolve_secret() db.session.execute(text("UPDATE public.identity_idp_config_versions SET status='superseded' WHERE provider_uid=CAST(:p AS uuid) AND status='active'"), {"p": provider_uid}) db.session.execute(text("UPDATE public.identity_idp_config_versions SET status='active',activated_at=CURRENT_TIMESTAMP WHERE provider_uid=CAST(:p AS uuid) AND version=:v"), {"p": provider_uid, "v": version}) db.session.execute(text("UPDATE public.enterprise_identity_links SET token_version=token_version+1,updated_at=CURRENT_TIMESTAMP WHERE provider_uid=CAST(:p AS uuid)"), {"p": provider_uid}) db.session.execute(text(""" UPDATE public.identity_sessions SET status='revoked',revoke_reason='provider_configuration_changed',revoked_at=CURRENT_TIMESTAMP WHERE status='active' AND identity_link_uid IN (SELECT uid FROM public.enterprise_identity_links WHERE provider_uid=CAST(:p AS uuid)) """), {"p": provider_uid}) elif action == "disable": disabled = db.session.execute(text("UPDATE public.identity_idp_config_versions SET status='disabled',disabled_at=CURRENT_TIMESTAMP WHERE provider_uid=CAST(:p AS uuid) AND version=:v AND status='active' RETURNING provider_uid"), {"p": provider_uid, "v": version}).scalar_one_or_none() if not disabled: raise IdentityPolicyError("only the selected active IdP version can be disabled") db.session.execute(text("UPDATE public.enterprise_identity_links SET token_version=token_version+1,updated_at=CURRENT_TIMESTAMP WHERE provider_uid=CAST(:p AS uuid)"), {"p": provider_uid}) db.session.execute(text(""" UPDATE public.identity_sessions SET status='revoked',revoke_reason='provider_disabled',revoked_at=CURRENT_TIMESTAMP WHERE status='active' AND identity_link_uid IN (SELECT uid FROM public.enterprise_identity_links WHERE provider_uid=CAST(:p AS uuid)) """), {"p": provider_uid}) elif action == "connectivity": payload = OidcTransport().fetch_jwks(config) result = {"reachable": True, "key_count": len(payload.get("keys", []))} _audit("idp_connectivity", "success", actor_uid=g.current_user["id"], provider_uid=provider_uid, resource_uid=f"{provider_uid}:{version}", detail={"version": version, **result}) return jsonify(success(result)) else: raise IdentityPolicyError("unsupported IdP version action") db.session.commit() _audit(f"idp_{action}", "success", actor_uid=g.current_user["id"], provider_uid=provider_uid) return jsonify(success({"provider_uid": provider_uid, "version": version, "action": action})) @bp.route("/identity/claims/dry-run", methods=["POST"]) @require_permissions(IDENTITY_OPERATE) @_identity_input_guard def identity_claims_dry_run(): data = _body() mapping_version = _required_str(data, "mapping_version", maximum=120) claims = data.get("claims") if not isinstance(claims, dict): raise IdentityPolicyError("claims must be an object") mapped = _validate_mapping(mapping_version, data.get("group_rules")).map(claims) result = {"roles": mapped.roles, "business_domain_uids": mapped.business_domain_uids, "object_types": mapped.object_types, "environments": mapped.environments, "data_scopes": mapped.data_scopes, "mapping_version": mapped.mapping_version, "claims_digest": mapped.claims_digest} _audit("claims_dry_run", "success", actor_uid=g.current_user["id"], enterprise_subject=mapped.subject, detail={"mapping_version": mapped.mapping_version, "claims_digest": mapped.claims_digest, "role_count": len(mapped.roles), "scope_count": sum(len(result[key]) for key in ("business_domain_uids", "object_types", "environments", "data_scopes"))}) return jsonify(success(result)) @bp.route("/identity/directory/delta", methods=["POST"]) @require_permissions(IDENTITY_OPERATE) @_identity_input_guard def identity_directory_delta(): data = _body() event_type = _required_str(data, "event_type", maximum=20).upper() provider_uid = _uuid(data.get("provider_uid"), "provider_uid") params = { "provider_uid": provider_uid, "source": _required_str(data, "source", maximum=80), "source_event_id": _required_str(data, "source_event_id", maximum=300), "cursor": _required_str(data, "cursor", maximum=500), "cursor_sequence": _positive_int(data.get("cursor_sequence"), "cursor_sequence"), "event_type": event_type, } if event_type in {"JOINER", "MOVER", "RESTORE"}: row = _load_config(provider_uid) if not row: raise IdentityPolicyError("active identity provider not found") config = _config(row) mapped = ClaimMapper( version=config.mapping_version, group_rules=dict(row["config"]).get("group_rules", {}), ).map(data.get("claims", {})) params["subject"] = mapped.subject params["attributes"] = { "provider_uid": config.provider_uid, "username": mapped.username, "display_name": mapped.display_name, "department": mapped.department, "groups": mapped.groups, "roles": mapped.roles, "mapping_version": mapped.mapping_version, "claims_digest": mapped.claims_digest, "authorization_scope": { "business_domain_uids": mapped.business_domain_uids, "object_types": mapped.object_types, "environments": mapped.environments, "data_scopes": mapped.data_scopes, }, } elif event_type in {"LEAVER", "DISABLE"}: params["subject"] = _required_str(data, "subject", maximum=500) params["attributes"] = {"username": data.get("username"), "roles": []} elif event_type in {"DEPARTMENT", "GROUP"}: if not isinstance(data.get("attributes"), dict): raise IdentityPolicyError("organization attributes must be an object") params["attributes"] = data["attributes"] else: raise IdentityPolicyError("unsupported directory event type") result = DirectorySynchronizer(_repo(), SessionManager(_repo(), secret=current_app.config["SECRET_KEY"])).apply(**params) _audit("directory_delta", "success", actor_uid=g.current_user["id"], provider_uid=provider_uid, resource_uid=params["source_event_id"], enterprise_subject=str(params["subject"]) if params.get("subject") else None, detail={"event_type": event_type, "result": result}) return jsonify(success(result)) @bp.route("/identity/sessions", methods=["GET"]) @require_permissions(IDENTITY_OPERATE) @_identity_input_guard def identity_sessions(): subject = request.args.get("subject") provider_uid = request.args.get("provider_uid") if subject and not provider_uid: raise IdentityPolicyError("provider_uid is required with subject") if provider_uid: provider_uid = _uuid(provider_uid, "provider_uid") if subject and len(subject) > 500: raise IdentityPolicyError("subject is invalid") rows = _repo().list_sessions(provider_uid, subject) visible = ("uid", "user_uid", "identity_source", "status", "created_at", "last_seen_at", "expires_at") return jsonify(success([{key: row.get(key) for key in visible} for row in rows])) @bp.route("/identity/sessions//revoke", methods=["POST"]) @require_permissions(IDENTITY_OPERATE) @_identity_input_guard def identity_session_revoke(session_uid: str): session_uid = _uuid(session_uid, "session_uid") data = _body() session = _repo().get_session(session_uid) reason = _optional_str(data, "reason", "operator_revoke", maximum=120) SessionManager(_repo(), secret=current_app.config["SECRET_KEY"]).revoke(session_uid, reason=reason) _audit("identity_session_revoke", "success", actor_uid=g.current_user["id"], provider_uid=session.get("provider_uid") if session else None, resource_uid=session_uid, detail={"reason_supplied": bool(data.get("reason"))}) return jsonify(success({"session_uid": session_uid, "revoked": True})) @bp.route("/identity/sessions//risk-terminate", methods=["POST"]) @require_permissions(IDENTITY_OPERATE) @_identity_input_guard def identity_session_risk(session_uid: str): session_uid = _uuid(session_uid, "session_uid") data = _body() session = _repo().get_session(session_uid) reason = _optional_str(data, "reason", "risk_signal", maximum=300) SessionManager(_repo(), secret=current_app.config["SECRET_KEY"]).risk_revoke(session_uid, reason=reason) _audit("identity_session_risk_terminate", "success", actor_uid=g.current_user["id"], provider_uid=session.get("provider_uid") if session else None, resource_uid=session_uid, detail={"reason_supplied": bool(data.get("reason"))}) return jsonify(success({"session_uid": session_uid, "revoked": True, "risk": True})) @bp.route("/identity/emergency/requests", methods=["POST"]) @require_permissions(IDENTITY_OPERATE) @_identity_input_guard def identity_emergency_request(): data = _body() account_uid = _uuid(data.get("account_uid"), "account_uid") reason = _required_str(data, "reason", maximum=2000) expires_at = _timestamp(data.get("expires_at"), "expires_at") eligible = db.session.execute(text(""" SELECT EXISTS( SELECT 1 FROM public.users u JOIN public.user_roles ur ON ur.user_id=u.id JOIN public.roles r ON r.id=ur.role_id WHERE u.id=CAST(:account AS uuid) AND u.status='active' AND r.name='admin' AND NOT EXISTS(SELECT 1 FROM public.enterprise_identity_links l WHERE l.user_uid=u.id) ) """), {"account": account_uid}).scalar() record = EmergencyAccess(_repo()).request(requester_uid=g.current_user["id"], account_uid=account_uid, reason=reason, expires_at=expires_at, account_is_local_active_admin=bool(eligible)) _audit("emergency_request", "success", actor_uid=g.current_user["id"], resource_uid=record["uid"], detail={"account_uid": record["account_uid"], "expires_at": record["expires_at"].isoformat()}) return jsonify(success(record)), 201 @bp.route("/identity/emergency//", methods=["POST"]) @require_permissions(IDENTITY_MANAGE) @_identity_input_guard def identity_emergency_action(request_uid: str, action: str): request_uid = _uuid(request_uid, "request_uid") service = EmergencyAccess(_repo()) data = _body() if action == "approve": result = service.approve(request_uid, approver_uid=g.current_user["id"]) elif action == "activate": result = service.activate(request_uid) elif action == "review": result = service.review(request_uid, reviewer_uid=g.current_user["id"], outcome=_required_str(data, "outcome", maximum=2000)) elif action == "close": result = service.close(request_uid, actor_uid=g.current_user["id"]) else: raise IdentityPolicyError("unsupported emergency action") _audit(f"emergency_{action}", "warning" if action == "activate" else "success", actor_uid=g.current_user["id"], resource_uid=request_uid) return jsonify(success(result)) @bp.route("/identity/sessions/revoke-all", methods=["POST"]) @require_permissions(IDENTITY_OPERATE) @_identity_input_guard def identity_session_revoke_all(): data = _body() provider_uid = _uuid(data.get("provider_uid"), "provider_uid") subject = _required_str(data, "subject", maximum=500) reason = _optional_str(data, "reason", "operator_revoke_all", maximum=120) SessionManager(_repo(), secret=current_app.config["SECRET_KEY"]).revoke_subject( provider_uid, subject, reason=reason ) _audit("identity_session_revoke_all", "success", actor_uid=g.current_user["id"], provider_uid=provider_uid, enterprise_subject=subject, detail={"reason_supplied": bool(data.get("reason"))}) return jsonify(success({"revoked": True})) @bp.route("/identity/audit", methods=["GET"]) @require_permissions(IDENTITY_MANAGE) @_identity_input_guard def identity_audit(): rows = db.session.execute(text("SELECT uid::text,event_type,outcome,actor_uid::text,provider_uid::text,resource_type,resource_uid,safe_detail,created_at FROM public.identity_audit_events ORDER BY created_at DESC LIMIT :limit"), {"limit": _positive_int(request.args.get("limit", 100), "limit", maximum=500)}).mappings().all() return jsonify(success([dict(row) for row in rows]))