auth.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329
  1. """PostgreSQL-backed authentication using Argon2id and short-lived bearer tokens."""
  2. from __future__ import annotations
  3. import hashlib
  4. import json
  5. import logging
  6. import threading
  7. import time
  8. from collections import defaultdict, deque
  9. from typing import Any
  10. from argon2 import PasswordHasher
  11. from argon2.exceptions import InvalidHashError, VerificationError, VerifyMismatchError
  12. from flask import current_app
  13. from sqlalchemy import text
  14. from sqlalchemy.exc import ProgrammingError, SQLAlchemyError
  15. from app import db
  16. from app.core.common.identifiers import new_governance_uid
  17. from app.core.system.tokens import TokenError, decode_access_token, issue_access_token
  18. logger = logging.getLogger(__name__)
  19. _password_hasher = PasswordHasher()
  20. _failure_lock = threading.Lock()
  21. _login_failures: dict[str, deque[float]] = defaultdict(deque)
  22. _FAILURE_WINDOW_SECONDS = 300
  23. _FAILURE_LIMIT = 5
  24. def hash_password(password: str) -> str:
  25. return _password_hasher.hash(password)
  26. def verify_password(encoded: str, password: str) -> bool:
  27. try:
  28. return bool(_password_hasher.verify(encoded, password))
  29. except (VerifyMismatchError, VerificationError, InvalidHashError):
  30. return False
  31. def _rate_key(username: str, ip_address: str | None) -> str:
  32. return f"{username.strip().lower()}|{ip_address or '-'}"
  33. def _check_rate_limit(key: str, *, now: float | None = None) -> None:
  34. now = now or time.time()
  35. with _failure_lock:
  36. failures = _login_failures[key]
  37. while failures and failures[0] <= now - _FAILURE_WINDOW_SECONDS:
  38. failures.popleft()
  39. if len(failures) >= _FAILURE_LIMIT:
  40. raise PermissionError("too many login attempts")
  41. def _record_failure(key: str, *, now: float | None = None) -> None:
  42. with _failure_lock:
  43. _login_failures[key].append(now or time.time())
  44. def _clear_failures(key: str) -> None:
  45. with _failure_lock:
  46. _login_failures.pop(key, None)
  47. def _audit(
  48. *,
  49. user_id: str | None,
  50. username: str,
  51. event_type: str,
  52. success: bool,
  53. ip_address: str | None,
  54. user_agent: str | None,
  55. detail: str | None = None,
  56. ) -> None:
  57. db.session.execute(
  58. text(
  59. "INSERT INTO public.auth_audit_events "
  60. "(user_id, username, event_type, success, ip_address, user_agent, detail) "
  61. "VALUES (CAST(:user_id AS uuid), :username, :event_type, :success, "
  62. ":ip_address, :user_agent, :detail)"
  63. ),
  64. {
  65. "user_id": user_id,
  66. "username": username[:64],
  67. "event_type": event_type,
  68. "success": success,
  69. "ip_address": (ip_address or "")[:64] or None,
  70. "user_agent": (user_agent or "")[:300] or None,
  71. "detail": (detail or "")[:500] or None,
  72. },
  73. )
  74. try:
  75. with db.session.begin_nested():
  76. db.session.execute(
  77. text(
  78. "INSERT INTO public.identity_audit_events "
  79. "(uid,event_type,outcome,actor_uid,enterprise_subject_digest,resource_type,resource_uid,safe_detail) "
  80. "VALUES (CAST(:uid AS uuid),:event_type,:outcome,CAST(:actor AS uuid),:subject_digest,"
  81. "'authentication',:resource_uid,CAST(:detail AS jsonb))"
  82. ),
  83. {
  84. "uid": new_governance_uid(),
  85. "event_type": event_type,
  86. "outcome": "success" if success else "failure",
  87. "actor": user_id,
  88. "subject_digest": hashlib.sha256(username.strip().lower().encode()).hexdigest(),
  89. "resource_uid": user_id,
  90. "detail": json.dumps({"ip_address": (ip_address or "")[:64] or None,
  91. "user_agent": (user_agent or "")[:300] or None,
  92. "reason": (detail or "")[:160] or None}),
  93. },
  94. )
  95. except ProgrammingError:
  96. # Rolling deployment before migration 470 retains the legacy audit row.
  97. pass
  98. def authenticate_user(
  99. username: str,
  100. password: str,
  101. *,
  102. ip_address: str | None = None,
  103. user_agent: str | None = None,
  104. ) -> dict[str, Any] | None:
  105. key = _rate_key(username, ip_address)
  106. _check_rate_limit(key)
  107. row = db.session.execute(
  108. text(
  109. "SELECT u.id::text, u.username, u.display_name, u.password_hash, "
  110. "u.status, COALESCE(array_agg(r.name ORDER BY r.name) "
  111. "FILTER (WHERE r.name IS NOT NULL), ARRAY[]::varchar[]) AS roles "
  112. "FROM public.users u "
  113. "LEFT JOIN public.user_roles ur ON ur.user_id = u.id "
  114. "LEFT JOIN public.roles r ON r.id = ur.role_id "
  115. "WHERE lower(u.username) = lower(:username) "
  116. "GROUP BY u.id"
  117. ),
  118. {"username": username.strip()},
  119. ).one_or_none()
  120. valid = bool(row and row[4] == "active" and verify_password(row[3], password))
  121. if not valid:
  122. _record_failure(key)
  123. _audit(
  124. user_id=row[0] if row else None,
  125. username=username,
  126. event_type="login",
  127. success=False,
  128. ip_address=ip_address,
  129. user_agent=user_agent,
  130. detail="invalid credentials or disabled account",
  131. )
  132. db.session.commit()
  133. return None
  134. # Once an enterprise IdP is active, local authentication is a controlled
  135. # break-glass path only. Missing identity tables mean pre-migration/no-IdP
  136. # compatibility; an active IdP always fails closed.
  137. try:
  138. active_provider = db.session.execute(
  139. text("SELECT EXISTS (SELECT 1 FROM public.identity_idp_config_versions WHERE status='active')")
  140. ).scalar()
  141. except SQLAlchemyError:
  142. db.session.rollback()
  143. logger.exception("enterprise identity state unavailable; local login denied")
  144. return None
  145. if active_provider:
  146. emergency_request = db.session.execute(
  147. text(
  148. "SELECT uid::text,expires_at FROM public.identity_emergency_requests "
  149. "WHERE emergency_account_uid=CAST(:id AS uuid) AND status='active' "
  150. "AND starts_at<=CURRENT_TIMESTAMP AND expires_at>CURRENT_TIMESTAMP "
  151. "AND EXISTS (SELECT 1 FROM public.user_roles ur JOIN public.roles r ON r.id=ur.role_id "
  152. "WHERE ur.user_id=CAST(:id AS uuid) AND r.name='admin') "
  153. "AND NOT EXISTS (SELECT 1 FROM public.enterprise_identity_links l WHERE l.user_uid=CAST(:id AS uuid)) "
  154. "ORDER BY expires_at ASC LIMIT 1"
  155. ),
  156. {"id": row[0]},
  157. ).one_or_none()
  158. if not emergency_request:
  159. _audit(user_id=row[0], username=row[1], event_type="local_login_blocked_by_idp", success=False,
  160. ip_address=ip_address, user_agent=user_agent, detail="active enterprise IdP requires approved emergency access")
  161. db.session.commit()
  162. return None
  163. db.session.execute(
  164. text("UPDATE public.users SET last_login_at = CURRENT_TIMESTAMP WHERE id = CAST(:id AS uuid)"),
  165. {"id": row[0]},
  166. )
  167. _audit(
  168. user_id=row[0],
  169. username=row[1],
  170. event_type="login",
  171. success=True,
  172. ip_address=ip_address,
  173. user_agent=user_agent,
  174. )
  175. db.session.commit()
  176. _clear_failures(key)
  177. roles = list(row[5])
  178. token = issue_access_token(
  179. user_id=row[0], roles=roles, secret=current_app.config["SECRET_KEY"]
  180. )
  181. if active_provider:
  182. from app.core.system.identity_repository import PostgresIdentityRepository
  183. from app.core.system.identity_sessions import SessionManager
  184. credentials = SessionManager(
  185. PostgresIdentityRepository(db.session), secret=current_app.config["SECRET_KEY"]
  186. ).create(
  187. subject=row[0],
  188. user_uid=row[0],
  189. roles=roles,
  190. identity_source="emergency",
  191. emergency_request_uid=emergency_request[0],
  192. emergency_expires_at=emergency_request[1],
  193. )
  194. token = credentials.access_token
  195. PostgresIdentityRepository(db.session).audit(
  196. {
  197. "event_type": "emergency_local_login",
  198. "outcome": "warning",
  199. "actor_uid": row[0],
  200. "resource_type": "emergency_access",
  201. "resource_uid": row[0],
  202. "safe_detail": {"alert_required": True},
  203. }
  204. )
  205. return {
  206. "id": row[0],
  207. "username": row[1],
  208. "display_name": row[2],
  209. "roles": roles,
  210. "token": token,
  211. }
  212. def load_identity(user_id: str) -> dict[str, Any] | None:
  213. row = db.session.execute(
  214. text(
  215. "SELECT u.id::text, u.username, u.display_name, u.status, "
  216. "COALESCE(array_agg(r.name ORDER BY r.name) FILTER "
  217. "(WHERE r.name IS NOT NULL), ARRAY[]::varchar[]) "
  218. "FROM public.users u "
  219. "LEFT JOIN public.user_roles ur ON ur.user_id = u.id "
  220. "LEFT JOIN public.roles r ON r.id = ur.role_id "
  221. "WHERE u.id = CAST(:id AS uuid) GROUP BY u.id"
  222. ),
  223. {"id": user_id},
  224. ).one_or_none()
  225. if not row or row[3] != "active":
  226. return None
  227. return {"id": row[0], "username": row[1], "display_name": row[2], "roles": list(row[4])}
  228. def load_identity_from_token(token: str, *, secret: str) -> dict[str, Any] | None:
  229. try:
  230. claims = decode_access_token(token, secret=secret)
  231. except TokenError:
  232. return None
  233. try:
  234. active_provider = bool(db.session.execute(text(
  235. "SELECT EXISTS (SELECT 1 FROM public.identity_idp_config_versions WHERE status='active')"
  236. )).scalar())
  237. except SQLAlchemyError:
  238. db.session.rollback()
  239. logger.exception("enterprise identity state unavailable; bearer token denied")
  240. return None
  241. if active_provider and not claims.get("sid"):
  242. return None
  243. if claims.get("sid"):
  244. try:
  245. session = db.session.execute(
  246. text("""
  247. SELECT s.status,s.token_version,s.user_uid::text,s.identity_source,
  248. CASE WHEN s.identity_source='emergency' THEN
  249. (e.status='active' AND e.starts_at<=CURRENT_TIMESTAMP AND e.expires_at>CURRENT_TIMESTAMP
  250. AND u.status='active'
  251. AND EXISTS (SELECT 1 FROM public.user_roles ur JOIN public.roles r ON r.id=ur.role_id
  252. WHERE ur.user_id=u.id AND r.name='admin')
  253. AND NOT EXISTS (SELECT 1 FROM public.enterprise_identity_links l WHERE l.user_uid=u.id))
  254. ELSE TRUE END emergency_valid
  255. FROM public.identity_sessions s LEFT JOIN public.identity_emergency_requests e
  256. ON e.uid=s.emergency_request_uid LEFT JOIN public.users u ON u.id=s.user_uid
  257. WHERE s.uid=CAST(:sid AS uuid)
  258. """),
  259. {"sid": claims["sid"]},
  260. ).one_or_none()
  261. except SQLAlchemyError:
  262. db.session.rollback()
  263. return None
  264. if session and session[3] == "emergency" and not session[4]:
  265. db.session.execute(text("UPDATE public.identity_sessions SET status='revoked',revoke_reason='emergency_window_closed',revoked_at=CURRENT_TIMESTAMP WHERE uid=CAST(:sid AS uuid)"), {"sid": claims["sid"]})
  266. db.session.commit()
  267. return None
  268. if (not session or session[0] != "active" or session[2] != claims["sub"]
  269. or int(session[1]) != int(claims.get("token_version", -1))):
  270. return None
  271. identity = load_identity(claims["sub"])
  272. if identity and claims.get("sid"):
  273. identity.update(sid=claims["sid"], token_version=claims["token_version"],
  274. identity_source=claims.get("identity_source", "oidc"))
  275. return identity
  276. # Compatibility aliases retained only for imports during the route migration.
  277. def login_user(username: str, password: str):
  278. result = authenticate_user(username, password)
  279. return (True, result) if result else (False, "用户名或密码错误")
  280. def get_user_by_username(username: str):
  281. row = db.session.execute(
  282. text("SELECT id::text FROM public.users WHERE lower(username) = lower(:username)"),
  283. {"username": username},
  284. ).scalar_one_or_none()
  285. return load_identity(row) if row else None
  286. def init_db() -> bool:
  287. logger.warning("init_db is deprecated; use Alembic migrations")
  288. return True
  289. def require_auth(view):
  290. from app.core.system.permissions import READ_GOVERNANCE, require_permissions
  291. return require_permissions(READ_GOVERNANCE)(view)