| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648 |
- """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/<provider_uid>/<int:version>/<action>", 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/<session_uid>/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/<session_uid>/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/<request_uid>/<action>", 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]))
|