enterprise_identity.py 34 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648
  1. """Enterprise identity APIs. Public routes are explicitly enumerated in policy."""
  2. from __future__ import annotations
  3. import hashlib
  4. import json
  5. import secrets
  6. import uuid
  7. from datetime import UTC, datetime, timedelta
  8. from functools import wraps
  9. from urllib.parse import urlsplit
  10. from flask import current_app, g, jsonify, redirect, request
  11. from itsdangerous import BadSignature, URLSafeTimedSerializer
  12. from sqlalchemy import text
  13. from sqlalchemy.exc import IntegrityError, StatementError
  14. from app import db
  15. from app.api.system import bp
  16. from app.core.common.identifiers import new_governance_uid
  17. from app.core.system.enterprise_identity import (
  18. ClaimMapper,
  19. DirectorySynchronizer,
  20. EmergencyAccess,
  21. IdentityPolicyError,
  22. IdentityUpstreamError,
  23. IdpConfig,
  24. sanitize_audit_detail,
  25. )
  26. from app.core.system.identity_repository import PostgresIdentityRepository
  27. from app.core.system.identity_sessions import SessionManager
  28. from app.core.system.oidc import OidcClient, OidcTransport
  29. from app.core.system.permissions import (
  30. IDENTITY_MANAGE,
  31. IDENTITY_OPERATE,
  32. IDENTITY_READ,
  33. require_permissions,
  34. )
  35. from app.core.system.tokens import issue_access_token
  36. from app.models.result import failed, success
  37. def _body() -> dict:
  38. value = request.get_json(silent=True)
  39. if not isinstance(value, dict):
  40. raise IdentityPolicyError("JSON object body is required")
  41. return value
  42. def _required_str(data: dict, key: str, *, maximum: int = 500) -> str:
  43. value = data.get(key)
  44. if not isinstance(value, str) or not value.strip() or len(value.strip()) > maximum:
  45. raise IdentityPolicyError(f"{key} is required or invalid")
  46. return value.strip()
  47. def _optional_str(data: dict, key: str, default: str, *, maximum: int) -> str:
  48. if key not in data:
  49. return default
  50. return _required_str(data, key, maximum=maximum)
  51. def _uuid(value, field: str) -> str:
  52. try:
  53. return str(uuid.UUID(str(value)))
  54. except (ValueError, TypeError, AttributeError) as exc:
  55. raise IdentityPolicyError(f"{field} must be a UUID") from exc
  56. def _positive_int(value, field: str, *, maximum: int = 2_147_483_647) -> int:
  57. if isinstance(value, bool):
  58. raise IdentityPolicyError(f"{field} must be a positive integer")
  59. try:
  60. result = int(value)
  61. except (ValueError, TypeError) as exc:
  62. raise IdentityPolicyError(f"{field} must be a positive integer") from exc
  63. if result < 1 or result > maximum:
  64. raise IdentityPolicyError(f"{field} must be a positive integer")
  65. return result
  66. def _timestamp(value, field: str) -> datetime:
  67. if not isinstance(value, str):
  68. raise IdentityPolicyError(f"{field} must be an ISO-8601 timestamp")
  69. try:
  70. parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
  71. except ValueError as exc:
  72. raise IdentityPolicyError(f"{field} must be an ISO-8601 timestamp") from exc
  73. if parsed.tzinfo is None or parsed.utcoffset() is None:
  74. raise IdentityPolicyError(f"{field} must include a timezone")
  75. return parsed
  76. def _string_array(data: dict, key: str, *, default=None) -> tuple[str, ...]:
  77. value = data.get(key, default)
  78. if not isinstance(value, list) or not value or any(not isinstance(item, str) or not item.strip() for item in value):
  79. raise IdentityPolicyError(f"{key} must be a non-empty string array")
  80. return tuple(item.strip() for item in value)
  81. def _validate_mapping(version: str, group_rules) -> ClaimMapper:
  82. if not isinstance(group_rules, dict):
  83. raise IdentityPolicyError("group_rules must be an object")
  84. return ClaimMapper(version=version, group_rules=group_rules)
  85. def _identity_input_guard(view):
  86. @wraps(view)
  87. def guarded(*args, **kwargs):
  88. try:
  89. return view(*args, **kwargs)
  90. except IdentityPolicyError:
  91. raise
  92. except (KeyError, TypeError, ValueError, IntegrityError, StatementError) as exc:
  93. db.session.rollback()
  94. raise IdentityPolicyError("invalid enterprise identity request") from exc
  95. return guarded
  96. def _repo() -> PostgresIdentityRepository:
  97. return PostgresIdentityRepository(db.session)
  98. def _config(row) -> IdpConfig:
  99. return IdpConfig(provider_uid=str(row["provider_uid"]), version=row["version"], issuer=row["issuer"],
  100. client_id=row["client_id"], secret_ref=row["secret_ref"],
  101. authorization_endpoint=row["authorization_endpoint"], token_endpoint=row["token_endpoint"],
  102. jwks_uri=row["jwks_uri"], redirect_uris=tuple(row["redirect_uris"]),
  103. post_login_redirect_uris=tuple(dict(row["config"]).get("post_login_redirect_uris", ())),
  104. algorithms=tuple(row["algorithms"]), mapping_version=row["mapping_version"], status=row["status"])
  105. def _load_config(provider_uid: str, *, version: int | None = None, active: bool = True):
  106. clause = "AND status='active'" if active else ""
  107. version_clause = "AND version=:version" if version is not None else ""
  108. return db.session.execute(text(f"""
  109. SELECT provider_uid::text,version,status,issuer,client_id,secret_ref,authorization_endpoint,
  110. token_endpoint,jwks_uri,redirect_uris,algorithms,mapping_version,config
  111. FROM public.identity_idp_config_versions WHERE provider_uid=CAST(:provider AS uuid) {version_clause} {clause}
  112. ORDER BY version DESC LIMIT 1
  113. """), {"provider": provider_uid, "version": version}).mappings().one_or_none()
  114. def _audit(event_type: str, outcome: str, *, actor_uid: str | None = None,
  115. provider_uid: str | None = None, resource_uid: str | None = None,
  116. enterprise_subject: str | None = None, detail: dict | None = None) -> None:
  117. _repo().audit({"event_type": event_type, "outcome": outcome, "actor_uid": actor_uid,
  118. "provider_uid": provider_uid, "resource_type": "enterprise_identity",
  119. "resource_uid": resource_uid,
  120. "subject_digest": (hashlib.sha256(enterprise_subject.encode()).hexdigest()
  121. if enterprise_subject else None),
  122. "safe_detail": sanitize_audit_detail(detail or {})})
  123. def _cookie_secure() -> bool:
  124. return bool(current_app.config.get("IDENTITY_COOKIE_SECURE", True))
  125. def _origin(uri: str) -> str:
  126. parts = urlsplit(uri)
  127. return f"{parts.scheme}://{parts.netloc}"
  128. def _browser_origin_allowed(origin: str) -> bool:
  129. rows = db.session.execute(text("SELECT config FROM public.identity_idp_config_versions WHERE status='active'")).scalars().all()
  130. allowed = {_origin(uri) for config in rows for uri in dict(config).get("post_login_redirect_uris", ())}
  131. allowed.update(current_app.config.get("IDENTITY_ALLOWED_TEST_ORIGINS", ()))
  132. return origin in allowed
  133. @bp.errorhandler(IdentityPolicyError)
  134. def identity_policy_error(exc):
  135. try:
  136. db.session.rollback()
  137. _audit("identity_policy_denied", "denied", actor_uid=getattr(g, "current_user", {}).get("id"),
  138. detail={"path": request.path, "reason": str(exc)[:160]})
  139. except Exception: # pragma: no cover - failure reporting must not expose internals
  140. db.session.rollback()
  141. return jsonify(failed(str(exc), code=400)), 400
  142. @bp.errorhandler(IdentityUpstreamError)
  143. def identity_upstream_error(exc):
  144. db.session.rollback()
  145. try:
  146. _audit("identity_upstream_failure", "failure", actor_uid=getattr(g, "current_user", {}).get("id"),
  147. detail={"path": request.path})
  148. except Exception: # pragma: no cover
  149. db.session.rollback()
  150. return jsonify(failed("enterprise identity provider is unavailable", code=502)), 502
  151. @bp.route("/identity/providers", methods=["GET"])
  152. def identity_providers():
  153. rows = db.session.execute(text("""
  154. SELECT provider_uid::text,version,issuer,config->>'display_name' display_name
  155. FROM public.identity_idp_config_versions WHERE status='active' ORDER BY provider_uid
  156. """)).mappings().all()
  157. return jsonify(success([{"provider_uid": row["provider_uid"], "version": row["version"],
  158. "issuer": row["issuer"], "display_name": row["display_name"] or "企业单点登录"} for row in rows]))
  159. @bp.route("/identity/authorize", methods=["POST"])
  160. @_identity_input_guard
  161. def identity_authorize():
  162. data = _body()
  163. provider_uid = _uuid(data.get("provider_uid"), "provider_uid")
  164. row = _load_config(provider_uid)
  165. if not row:
  166. raise IdentityPolicyError("active identity provider not found")
  167. config = _config(row)
  168. callback_uri = _required_str(data, "redirect_uri", maximum=2048)
  169. post_login_uri = _required_str(data, "post_login_uri", maximum=2048)
  170. if post_login_uri not in config.post_login_redirect_uris:
  171. raise IdentityPolicyError("post-login redirect URI is not allowlisted")
  172. started = OidcClient(_repo()).begin(config, callback_uri)
  173. serializer = URLSafeTimedSerializer(current_app.config["SECRET_KEY"], salt="oidc-flow")
  174. response = jsonify(success({"authorization_url": started.url, "provider_uid": config.provider_uid}))
  175. response.set_cookie("dataops_oidc_flow", serializer.dumps({"state": started.state, "verifier": started.code_verifier,
  176. "provider_uid": config.provider_uid,
  177. "redirect_uri": callback_uri,
  178. "post_login_uri": post_login_uri}),
  179. httponly=True, secure=_cookie_secure(), samesite="Lax", max_age=300, path="/api/system/identity/callback")
  180. _audit("oidc_authorize", "success", provider_uid=config.provider_uid, detail={"redirect_uri": data["redirect_uri"]})
  181. return response
  182. @bp.route("/identity/callback", methods=["GET"])
  183. @_identity_input_guard
  184. def identity_callback():
  185. serializer = URLSafeTimedSerializer(current_app.config["SECRET_KEY"], salt="oidc-flow")
  186. try:
  187. flow = serializer.loads(request.cookies.get("dataops_oidc_flow", ""), max_age=300)
  188. except BadSignature as exc:
  189. raise IdentityPolicyError("OIDC browser flow binding is invalid or expired") from exc
  190. if not isinstance(flow, dict) or request.args.get("state") != flow.get("state"):
  191. raise IdentityPolicyError("OIDC state mismatch")
  192. provider_uid = _uuid(flow.get("provider_uid"), "provider_uid")
  193. row = _load_config(provider_uid)
  194. if not row:
  195. raise IdentityPolicyError("active identity provider not found")
  196. row = dict(row)
  197. config = _config(row)
  198. config.validate()
  199. callback_uri = _required_str(flow, "redirect_uri", maximum=2048)
  200. verifier = _required_str(flow, "verifier", maximum=512)
  201. state = _required_str(flow, "state", maximum=512)
  202. post_login_uri = _required_str(flow, "post_login_uri", maximum=2048)
  203. code = request.args.get("code")
  204. if not isinstance(code, str) or not code or len(code) > 4096:
  205. raise IdentityPolicyError("OIDC authorization code is required")
  206. # Do not keep a database transaction open while calling the external IdP.
  207. db.session.rollback()
  208. transport = OidcTransport()
  209. id_token = transport.exchange_code(config, code=code, redirect_uri=callback_uri, verifier=verifier)
  210. jwks = transport.fetch_jwks(config)
  211. repo = _repo()
  212. try:
  213. claims = OidcClient(repo).verify_callback(config, state=state, redirect_uri=callback_uri,
  214. code_verifier=verifier, id_token=id_token, jwks=jwks,
  215. commit=False)
  216. mapped = ClaimMapper(version=config.mapping_version,
  217. group_rules=dict(row["config"]).get("group_rules", {})).map(claims)
  218. link = repo.get_identity(config.provider_uid, mapped.subject, for_update=True)
  219. if link and link["status"] != "active":
  220. raise IdentityPolicyError("enterprise identity is disabled")
  221. scope = {"business_domain_uids": mapped.business_domain_uids, "object_types": mapped.object_types,
  222. "environments": mapped.environments, "data_scopes": mapped.data_scopes}
  223. changed = bool(link and (sorted(link["roles"]) != sorted(mapped.roles)
  224. or dict(link["authorization_scope"]) != json.loads(json.dumps(scope))))
  225. user_uid = link["user_uid"] if link else new_governance_uid()
  226. token_version = (int(link["token_version"]) + 1 if changed else int(link["token_version"])) if link else 1
  227. repo.put_identity(config.provider_uid, mapped.subject, {
  228. "provider_uid": config.provider_uid, "user_uid": user_uid, "username": mapped.username,
  229. "display_name": mapped.display_name, "department": mapped.department, "groups": list(mapped.groups),
  230. "roles": list(mapped.roles), "authorization_scope": scope, "mapping_version": mapped.mapping_version,
  231. "claims_digest": mapped.claims_digest, "status": "active", "token_version": token_version,
  232. }, commit=False)
  233. if changed:
  234. repo.revoke_subject_sessions(config.provider_uid, mapped.subject,
  235. reason="claims_mapping_changed", commit=False)
  236. credentials = SessionManager(repo, secret=current_app.config["SECRET_KEY"]).create(
  237. subject=mapped.subject, user_uid=user_uid, roles=list(mapped.roles), identity_source="oidc",
  238. token_version=token_version, provider_uid=config.provider_uid, commit=False)
  239. exchange = secrets.token_urlsafe(40)
  240. repo.insert_exchange_code({"uid": new_governance_uid(),
  241. "code_hash": hashlib.sha256(exchange.encode()).hexdigest(),
  242. "session_uid": credentials.session_uid, "redirect_uri": post_login_uri,
  243. "expires_at": datetime.now(UTC) + timedelta(minutes=2)}, commit=False)
  244. db.session.commit()
  245. except IdentityPolicyError:
  246. db.session.rollback()
  247. raise
  248. except Exception as exc:
  249. db.session.rollback()
  250. raise IdentityPolicyError("OIDC callback persistence failed") from exc
  251. _audit("oidc_login", "success", provider_uid=config.provider_uid, resource_uid=user_uid,
  252. detail={"mapping_version": mapped.mapping_version, "claims_digest": mapped.claims_digest})
  253. response = redirect(f"{flow['post_login_uri']}?exchange_code={exchange}", code=302)
  254. response.set_cookie("dataops_refresh", credentials.refresh_token, httponly=True, secure=_cookie_secure(), samesite="Strict",
  255. max_age=8 * 3600, path="/api/system/identity")
  256. response.delete_cookie("dataops_oidc_flow", path="/api/system/identity/callback")
  257. return response
  258. @bp.route("/identity/exchange", methods=["POST"])
  259. @_identity_input_guard
  260. def identity_exchange():
  261. data = _body()
  262. exchange_code = _required_str(data, "exchange_code", maximum=1024)
  263. redirect_uri = _required_str(data, "redirect_uri", maximum=2048)
  264. digest = hashlib.sha256(exchange_code.encode()).hexdigest()
  265. row = db.session.execute(text("""
  266. UPDATE public.identity_exchange_codes SET status='consumed',consumed_at=CURRENT_TIMESTAMP
  267. WHERE code_hash=:digest AND status='pending' AND expires_at>CURRENT_TIMESTAMP
  268. AND redirect_uri=:redirect RETURNING session_uid::text,redirect_uri
  269. """), {"digest": digest, "redirect": redirect_uri}).mappings().one_or_none()
  270. if not row or request.headers.get("Origin") != _origin(row["redirect_uri"]):
  271. db.session.rollback()
  272. raise IdentityPolicyError("exchange code or browser origin is invalid")
  273. db.session.commit()
  274. session = _repo().get_session(row["session_uid"])
  275. if not session or session["status"] != "active":
  276. raise IdentityPolicyError("identity session is unavailable")
  277. token = issue_access_token(user_id=session["user_uid"], roles=list(session["roles"]), secret=current_app.config["SECRET_KEY"],
  278. session_uid=session["uid"], token_version=session["token_version"], identity_source=session["identity_source"])
  279. return jsonify(success({"token": token, "identity_source": session["identity_source"]}))
  280. @bp.route("/identity/refresh", methods=["POST"])
  281. def identity_refresh():
  282. origin = request.headers.get("Origin", "")
  283. if not _browser_origin_allowed(origin):
  284. raise IdentityPolicyError("browser origin is not allowlisted")
  285. credentials = SessionManager(_repo(), secret=current_app.config["SECRET_KEY"]).refresh(request.cookies.get("dataops_refresh", ""))
  286. session = _repo().get_session(credentials.session_uid)
  287. _audit("identity_refresh", "success", actor_uid=session["user_uid"] if session else None,
  288. provider_uid=session.get("provider_uid") if session else None,
  289. resource_uid=credentials.session_uid,
  290. detail={"rotated_from_uid": session.get("rotated_from_uid") if session else None})
  291. response = jsonify(success({"token": credentials.access_token}))
  292. response.set_cookie("dataops_refresh", credentials.refresh_token, httponly=True, secure=_cookie_secure(), samesite="Strict",
  293. max_age=8 * 3600, path="/api/system/identity")
  294. return response
  295. @bp.route("/identity/logout", methods=["POST"])
  296. @require_permissions(IDENTITY_READ)
  297. def identity_logout():
  298. session_uid = g.current_user.get("sid")
  299. session = _repo().get_session(session_uid) if session_uid else None
  300. if session_uid:
  301. SessionManager(_repo(), secret=current_app.config["SECRET_KEY"]).revoke(session_uid)
  302. _audit("identity_logout", "success", actor_uid=g.current_user.get("id"),
  303. provider_uid=session.get("provider_uid") if session else None, resource_uid=session_uid)
  304. response = jsonify(success({"revoked": True}))
  305. response.delete_cookie("dataops_refresh", path="/api/system/identity")
  306. return response
  307. @bp.route("/identity/idp-versions", methods=["GET"])
  308. @require_permissions(IDENTITY_MANAGE)
  309. def identity_idp_versions():
  310. 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()
  311. return jsonify(success([{**dict(row), "secret_ref": "env:DATAOPS_OIDC_***"} for row in rows]))
  312. @bp.route("/identity/idp-versions", methods=["POST"])
  313. @require_permissions(IDENTITY_MANAGE)
  314. @_identity_input_guard
  315. def identity_idp_version_create():
  316. data = _body()
  317. provider_uid = _uuid(data["provider_uid"], "provider_uid") if data.get("provider_uid") else new_governance_uid()
  318. mapping_version = _required_str(data, "mapping_version", maximum=120)
  319. group_rules = data.get("group_rules")
  320. _validate_mapping(mapping_version, group_rules)
  321. config = IdpConfig(provider_uid=provider_uid, version=_positive_int(data.get("version"), "version"),
  322. issuer=_required_str(data, "issuer", maximum=2048), client_id=_required_str(data, "client_id", maximum=300),
  323. secret_ref=_required_str(data, "secret_ref", maximum=160),
  324. authorization_endpoint=_required_str(data, "authorization_endpoint", maximum=2048),
  325. token_endpoint=_required_str(data, "token_endpoint", maximum=2048),
  326. jwks_uri=_required_str(data, "jwks_uri", maximum=2048), redirect_uris=_string_array(data, "redirect_uris"),
  327. post_login_redirect_uris=_string_array(data, "post_login_redirect_uris"),
  328. algorithms=_string_array(data, "algorithms", default=["RS256"]), mapping_version=mapping_version).validate()
  329. db.session.execute(text("""
  330. INSERT INTO public.identity_idp_config_versions(uid,provider_uid,version,status,issuer,client_id,secret_ref,authorization_endpoint,
  331. token_endpoint,jwks_uri,redirect_uris,algorithms,mapping_version,config,created_by)
  332. VALUES(CAST(:uid AS uuid),CAST(:provider AS uuid),:version,'draft',:issuer,:client,:secret,:auth,:token,:jwks,
  333. CAST(:redirects AS jsonb),CAST(:algorithms AS jsonb),:mapping,CAST(:config AS jsonb),CAST(:actor AS uuid))
  334. """), {"uid": new_governance_uid(), "provider": config.provider_uid, "version": config.version, "issuer": config.issuer,
  335. "client": config.client_id, "secret": config.secret_ref, "auth": config.authorization_endpoint,
  336. "token": config.token_endpoint, "jwks": config.jwks_uri, "redirects": json.dumps(config.redirect_uris),
  337. "algorithms": json.dumps(config.algorithms), "mapping": config.mapping_version,
  338. "config": json.dumps({"display_name": data.get("display_name"), "group_rules": group_rules,
  339. "post_login_redirect_uris": data.get("post_login_redirect_uris", [])}), "actor": g.current_user["id"]})
  340. db.session.commit()
  341. _audit("idp_version_created", "success", actor_uid=g.current_user["id"], provider_uid=config.provider_uid)
  342. return jsonify(success(config.public_dict())), 201
  343. @bp.route("/identity/idp-versions/<provider_uid>/<int:version>/<action>", methods=["POST"])
  344. @require_permissions(IDENTITY_MANAGE)
  345. @_identity_input_guard
  346. def identity_idp_action(provider_uid: str, version: int, action: str):
  347. provider_uid = _uuid(provider_uid, "provider_uid")
  348. version = _positive_int(version, "version")
  349. row = _load_config(provider_uid, version=version, active=False)
  350. if not row:
  351. raise IdentityPolicyError("IdP configuration version not found")
  352. config = _config(row)
  353. config.validate()
  354. _validate_mapping(config.mapping_version, dict(row["config"]).get("group_rules"))
  355. if action == "dry-run":
  356. result = {"valid": True, "secret_available": bool(config.resolve_secret())}
  357. _audit("idp_dry_run", "success", actor_uid=g.current_user["id"], provider_uid=provider_uid,
  358. resource_uid=f"{provider_uid}:{version}", detail={"version": version})
  359. return jsonify(success(result))
  360. if action in {"activate", "rotate"}:
  361. config.resolve_secret()
  362. 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})
  363. 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})
  364. 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})
  365. db.session.execute(text("""
  366. UPDATE public.identity_sessions SET status='revoked',revoke_reason='provider_configuration_changed',revoked_at=CURRENT_TIMESTAMP
  367. WHERE status='active' AND identity_link_uid IN
  368. (SELECT uid FROM public.enterprise_identity_links WHERE provider_uid=CAST(:p AS uuid))
  369. """), {"p": provider_uid})
  370. elif action == "disable":
  371. 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()
  372. if not disabled:
  373. raise IdentityPolicyError("only the selected active IdP version can be disabled")
  374. 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})
  375. db.session.execute(text("""
  376. UPDATE public.identity_sessions SET status='revoked',revoke_reason='provider_disabled',revoked_at=CURRENT_TIMESTAMP
  377. WHERE status='active' AND identity_link_uid IN
  378. (SELECT uid FROM public.enterprise_identity_links WHERE provider_uid=CAST(:p AS uuid))
  379. """), {"p": provider_uid})
  380. elif action == "connectivity":
  381. payload = OidcTransport().fetch_jwks(config)
  382. result = {"reachable": True, "key_count": len(payload.get("keys", []))}
  383. _audit("idp_connectivity", "success", actor_uid=g.current_user["id"], provider_uid=provider_uid,
  384. resource_uid=f"{provider_uid}:{version}", detail={"version": version, **result})
  385. return jsonify(success(result))
  386. else:
  387. raise IdentityPolicyError("unsupported IdP version action")
  388. db.session.commit()
  389. _audit(f"idp_{action}", "success", actor_uid=g.current_user["id"], provider_uid=provider_uid)
  390. return jsonify(success({"provider_uid": provider_uid, "version": version, "action": action}))
  391. @bp.route("/identity/claims/dry-run", methods=["POST"])
  392. @require_permissions(IDENTITY_OPERATE)
  393. @_identity_input_guard
  394. def identity_claims_dry_run():
  395. data = _body()
  396. mapping_version = _required_str(data, "mapping_version", maximum=120)
  397. claims = data.get("claims")
  398. if not isinstance(claims, dict):
  399. raise IdentityPolicyError("claims must be an object")
  400. mapped = _validate_mapping(mapping_version, data.get("group_rules")).map(claims)
  401. result = {"roles": mapped.roles, "business_domain_uids": mapped.business_domain_uids,
  402. "object_types": mapped.object_types, "environments": mapped.environments,
  403. "data_scopes": mapped.data_scopes, "mapping_version": mapped.mapping_version,
  404. "claims_digest": mapped.claims_digest}
  405. _audit("claims_dry_run", "success", actor_uid=g.current_user["id"],
  406. enterprise_subject=mapped.subject,
  407. detail={"mapping_version": mapped.mapping_version, "claims_digest": mapped.claims_digest,
  408. "role_count": len(mapped.roles), "scope_count": sum(len(result[key]) for key in
  409. ("business_domain_uids", "object_types", "environments", "data_scopes"))})
  410. return jsonify(success(result))
  411. @bp.route("/identity/directory/delta", methods=["POST"])
  412. @require_permissions(IDENTITY_OPERATE)
  413. @_identity_input_guard
  414. def identity_directory_delta():
  415. data = _body()
  416. event_type = _required_str(data, "event_type", maximum=20).upper()
  417. provider_uid = _uuid(data.get("provider_uid"), "provider_uid")
  418. params = {
  419. "provider_uid": provider_uid,
  420. "source": _required_str(data, "source", maximum=80),
  421. "source_event_id": _required_str(data, "source_event_id", maximum=300),
  422. "cursor": _required_str(data, "cursor", maximum=500),
  423. "cursor_sequence": _positive_int(data.get("cursor_sequence"), "cursor_sequence"),
  424. "event_type": event_type,
  425. }
  426. if event_type in {"JOINER", "MOVER", "RESTORE"}:
  427. row = _load_config(provider_uid)
  428. if not row:
  429. raise IdentityPolicyError("active identity provider not found")
  430. config = _config(row)
  431. mapped = ClaimMapper(
  432. version=config.mapping_version,
  433. group_rules=dict(row["config"]).get("group_rules", {}),
  434. ).map(data.get("claims", {}))
  435. params["subject"] = mapped.subject
  436. params["attributes"] = {
  437. "provider_uid": config.provider_uid,
  438. "username": mapped.username,
  439. "display_name": mapped.display_name,
  440. "department": mapped.department,
  441. "groups": mapped.groups,
  442. "roles": mapped.roles,
  443. "mapping_version": mapped.mapping_version,
  444. "claims_digest": mapped.claims_digest,
  445. "authorization_scope": {
  446. "business_domain_uids": mapped.business_domain_uids,
  447. "object_types": mapped.object_types,
  448. "environments": mapped.environments,
  449. "data_scopes": mapped.data_scopes,
  450. },
  451. }
  452. elif event_type in {"LEAVER", "DISABLE"}:
  453. params["subject"] = _required_str(data, "subject", maximum=500)
  454. params["attributes"] = {"username": data.get("username"), "roles": []}
  455. elif event_type in {"DEPARTMENT", "GROUP"}:
  456. if not isinstance(data.get("attributes"), dict):
  457. raise IdentityPolicyError("organization attributes must be an object")
  458. params["attributes"] = data["attributes"]
  459. else:
  460. raise IdentityPolicyError("unsupported directory event type")
  461. result = DirectorySynchronizer(_repo(), SessionManager(_repo(), secret=current_app.config["SECRET_KEY"])).apply(**params)
  462. _audit("directory_delta", "success", actor_uid=g.current_user["id"],
  463. provider_uid=provider_uid, resource_uid=params["source_event_id"],
  464. enterprise_subject=str(params["subject"]) if params.get("subject") else None,
  465. detail={"event_type": event_type, "result": result})
  466. return jsonify(success(result))
  467. @bp.route("/identity/sessions", methods=["GET"])
  468. @require_permissions(IDENTITY_OPERATE)
  469. @_identity_input_guard
  470. def identity_sessions():
  471. subject = request.args.get("subject")
  472. provider_uid = request.args.get("provider_uid")
  473. if subject and not provider_uid:
  474. raise IdentityPolicyError("provider_uid is required with subject")
  475. if provider_uid:
  476. provider_uid = _uuid(provider_uid, "provider_uid")
  477. if subject and len(subject) > 500:
  478. raise IdentityPolicyError("subject is invalid")
  479. rows = _repo().list_sessions(provider_uid, subject)
  480. visible = ("uid", "user_uid", "identity_source", "status", "created_at", "last_seen_at", "expires_at")
  481. return jsonify(success([{key: row.get(key) for key in visible} for row in rows]))
  482. @bp.route("/identity/sessions/<session_uid>/revoke", methods=["POST"])
  483. @require_permissions(IDENTITY_OPERATE)
  484. @_identity_input_guard
  485. def identity_session_revoke(session_uid: str):
  486. session_uid = _uuid(session_uid, "session_uid")
  487. data = _body()
  488. session = _repo().get_session(session_uid)
  489. reason = _optional_str(data, "reason", "operator_revoke", maximum=120)
  490. SessionManager(_repo(), secret=current_app.config["SECRET_KEY"]).revoke(session_uid, reason=reason)
  491. _audit("identity_session_revoke", "success", actor_uid=g.current_user["id"],
  492. provider_uid=session.get("provider_uid") if session else None, resource_uid=session_uid,
  493. detail={"reason_supplied": bool(data.get("reason"))})
  494. return jsonify(success({"session_uid": session_uid, "revoked": True}))
  495. @bp.route("/identity/sessions/<session_uid>/risk-terminate", methods=["POST"])
  496. @require_permissions(IDENTITY_OPERATE)
  497. @_identity_input_guard
  498. def identity_session_risk(session_uid: str):
  499. session_uid = _uuid(session_uid, "session_uid")
  500. data = _body()
  501. session = _repo().get_session(session_uid)
  502. reason = _optional_str(data, "reason", "risk_signal", maximum=300)
  503. SessionManager(_repo(), secret=current_app.config["SECRET_KEY"]).risk_revoke(session_uid, reason=reason)
  504. _audit("identity_session_risk_terminate", "success", actor_uid=g.current_user["id"],
  505. provider_uid=session.get("provider_uid") if session else None, resource_uid=session_uid,
  506. detail={"reason_supplied": bool(data.get("reason"))})
  507. return jsonify(success({"session_uid": session_uid, "revoked": True, "risk": True}))
  508. @bp.route("/identity/emergency/requests", methods=["POST"])
  509. @require_permissions(IDENTITY_OPERATE)
  510. @_identity_input_guard
  511. def identity_emergency_request():
  512. data = _body()
  513. account_uid = _uuid(data.get("account_uid"), "account_uid")
  514. reason = _required_str(data, "reason", maximum=2000)
  515. expires_at = _timestamp(data.get("expires_at"), "expires_at")
  516. eligible = db.session.execute(text("""
  517. SELECT EXISTS(
  518. SELECT 1 FROM public.users u JOIN public.user_roles ur ON ur.user_id=u.id
  519. JOIN public.roles r ON r.id=ur.role_id
  520. WHERE u.id=CAST(:account AS uuid) AND u.status='active' AND r.name='admin'
  521. AND NOT EXISTS(SELECT 1 FROM public.enterprise_identity_links l WHERE l.user_uid=u.id)
  522. )
  523. """), {"account": account_uid}).scalar()
  524. record = EmergencyAccess(_repo()).request(requester_uid=g.current_user["id"], account_uid=account_uid,
  525. reason=reason, expires_at=expires_at,
  526. account_is_local_active_admin=bool(eligible))
  527. _audit("emergency_request", "success", actor_uid=g.current_user["id"], resource_uid=record["uid"],
  528. detail={"account_uid": record["account_uid"], "expires_at": record["expires_at"].isoformat()})
  529. return jsonify(success(record)), 201
  530. @bp.route("/identity/emergency/<request_uid>/<action>", methods=["POST"])
  531. @require_permissions(IDENTITY_MANAGE)
  532. @_identity_input_guard
  533. def identity_emergency_action(request_uid: str, action: str):
  534. request_uid = _uuid(request_uid, "request_uid")
  535. service = EmergencyAccess(_repo())
  536. data = _body()
  537. if action == "approve":
  538. result = service.approve(request_uid, approver_uid=g.current_user["id"])
  539. elif action == "activate":
  540. result = service.activate(request_uid)
  541. elif action == "review":
  542. result = service.review(request_uid, reviewer_uid=g.current_user["id"],
  543. outcome=_required_str(data, "outcome", maximum=2000))
  544. elif action == "close":
  545. result = service.close(request_uid, actor_uid=g.current_user["id"])
  546. else:
  547. raise IdentityPolicyError("unsupported emergency action")
  548. _audit(f"emergency_{action}", "warning" if action == "activate" else "success", actor_uid=g.current_user["id"], resource_uid=request_uid)
  549. return jsonify(success(result))
  550. @bp.route("/identity/sessions/revoke-all", methods=["POST"])
  551. @require_permissions(IDENTITY_OPERATE)
  552. @_identity_input_guard
  553. def identity_session_revoke_all():
  554. data = _body()
  555. provider_uid = _uuid(data.get("provider_uid"), "provider_uid")
  556. subject = _required_str(data, "subject", maximum=500)
  557. reason = _optional_str(data, "reason", "operator_revoke_all", maximum=120)
  558. SessionManager(_repo(), secret=current_app.config["SECRET_KEY"]).revoke_subject(
  559. provider_uid, subject, reason=reason
  560. )
  561. _audit("identity_session_revoke_all", "success", actor_uid=g.current_user["id"],
  562. provider_uid=provider_uid, enterprise_subject=subject,
  563. detail={"reason_supplied": bool(data.get("reason"))})
  564. return jsonify(success({"revoked": True}))
  565. @bp.route("/identity/audit", methods=["GET"])
  566. @require_permissions(IDENTITY_MANAGE)
  567. @_identity_input_guard
  568. def identity_audit():
  569. 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"),
  570. {"limit": _positive_int(request.args.get("limit", 100), "limit", maximum=500)}).mappings().all()
  571. return jsonify(success([dict(row) for row in rows]))