identity.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379
  1. """Short-lived, one-time-issued connector machine credentials."""
  2. from __future__ import annotations
  3. import hashlib
  4. import hmac
  5. import secrets
  6. from collections.abc import Mapping
  7. from sqlalchemy import text
  8. from app.core.common.identifiers import new_governance_uid
  9. from app.core.connectors.errors import (
  10. ConnectorAuthenticationError,
  11. ConnectorConfigurationError,
  12. ConnectorPermissionError,
  13. )
  14. ALLOWED_OPERATIONS = frozenset(
  15. {
  16. "discover",
  17. "snapshot",
  18. "incremental",
  19. "lineage",
  20. "profile",
  21. "cancel",
  22. "resume",
  23. "evidence",
  24. }
  25. )
  26. ALLOWED_SCOPE_KEYS = frozenset(
  27. {"include_schemas", "exclude_schemas", "include_tables", "exclude_tables"}
  28. )
  29. def validate_machine_scope(scope):
  30. if not isinstance(scope, Mapping) or set(scope) - ALLOWED_SCOPE_KEYS:
  31. raise ConnectorConfigurationError("machine identity scope is invalid")
  32. normalized = {}
  33. for key, values in scope.items():
  34. if not isinstance(values, list) or any(
  35. not isinstance(value, str) or not value.strip() for value in values
  36. ):
  37. raise ConnectorConfigurationError(
  38. "machine identity scope values are invalid"
  39. )
  40. normalized[key] = tuple(value.strip() for value in values)
  41. return normalized
  42. def _token_hash(token):
  43. return hashlib.sha256(str(token).encode()).hexdigest()
  44. class ConnectorIdentityRepository:
  45. def __init__(self, session):
  46. self.session = session
  47. def create_principal(
  48. self,
  49. *,
  50. connector_id,
  51. connector_version,
  52. source_uid,
  53. business_domain_uid,
  54. environment,
  55. operations,
  56. scopes,
  57. actor_uid,
  58. source_binding_uid=None,
  59. source_binding_version=None,
  60. ):
  61. operations = tuple(sorted(set(operations)))
  62. if not operations or set(operations) - ALLOWED_OPERATIONS:
  63. raise ConnectorConfigurationError("machine identity operations are invalid")
  64. scopes = validate_machine_scope(scopes)
  65. if environment not in {"development", "staging", "production"}:
  66. raise ConnectorConfigurationError("machine identity environment is invalid")
  67. if not source_binding_uid or source_binding_version is None:
  68. raise ConnectorConfigurationError(
  69. "enterprise connector principal requires an approved source binding"
  70. )
  71. uid = new_governance_uid()
  72. self.session.execute(
  73. text("""
  74. INSERT INTO public.connector_principals
  75. (uid, connector_id, connector_version, source_uid, business_domain_uid, environment,
  76. allowed_operations, allowed_scopes, status, created_by,
  77. source_binding_uid,source_binding_version)
  78. VALUES (CAST(:uid AS uuid), :connector, :version, CAST(:source AS uuid), CAST(:domain AS uuid), :environment,
  79. CAST(:operations AS text[]), CAST(:scopes AS jsonb), 'active', CAST(:actor AS uuid),
  80. CAST(:binding AS uuid),:binding_version)
  81. """),
  82. {
  83. "uid": uid,
  84. "connector": connector_id,
  85. "version": connector_version,
  86. "source": source_uid,
  87. "domain": business_domain_uid,
  88. "environment": environment,
  89. "operations": list(operations),
  90. "scopes": __import__("json").dumps(scopes),
  91. "actor": actor_uid,
  92. "binding": source_binding_uid,
  93. "binding_version": source_binding_version,
  94. },
  95. )
  96. self._audit(uid, None, "principal_created", actor_uid, True)
  97. self.session.commit()
  98. return uid
  99. def issue(self, principal_uid, *, ttl_seconds, actor_uid, rotated_from=None):
  100. ttl = int(ttl_seconds)
  101. if ttl < 60 or ttl > 900:
  102. raise ConnectorConfigurationError(
  103. "credential TTL must be between 60 and 900 seconds"
  104. )
  105. self.session.execute(
  106. text("SELECT pg_advisory_xact_lock(hashtext(:key))"),
  107. {"key": f"connector-principal:{principal_uid}"},
  108. )
  109. principal = self.session.execute(
  110. text("""
  111. SELECT p.status
  112. FROM public.connector_principals p
  113. JOIN public.connector_source_bindings b
  114. ON b.uid=p.source_binding_uid
  115. AND b.binding_version=p.source_binding_version
  116. AND b.status='approved'
  117. AND b.connector_id=p.connector_id
  118. AND b.connector_version=p.connector_version
  119. AND b.source_uid=p.source_uid
  120. AND b.business_domain_uid=p.business_domain_uid
  121. AND b.environment=p.environment
  122. WHERE p.uid=CAST(:uid AS uuid)
  123. FOR UPDATE OF p
  124. """),
  125. {"uid": principal_uid},
  126. ).scalar_one_or_none()
  127. if principal != "active":
  128. raise ConnectorAuthenticationError(
  129. "connector principal or approved binding is not active"
  130. )
  131. token = "dopc_" + secrets.token_urlsafe(32)
  132. credential_uid = new_governance_uid()
  133. self.session.execute(
  134. text("""
  135. INSERT INTO public.connector_machine_credentials
  136. (uid, principal_uid, token_hash, status, expires_at, rotated_from_uid, issued_by)
  137. VALUES (CAST(:uid AS uuid), CAST(:principal AS uuid), :hash, 'active',
  138. CURRENT_TIMESTAMP + (:ttl * INTERVAL '1 second'), CAST(:rotated AS uuid), CAST(:actor AS uuid))
  139. """),
  140. {
  141. "uid": credential_uid,
  142. "principal": principal_uid,
  143. "hash": _token_hash(token),
  144. "ttl": ttl,
  145. "rotated": rotated_from,
  146. "actor": actor_uid,
  147. },
  148. )
  149. self._audit(principal_uid, credential_uid, "credential_issued", actor_uid, True)
  150. self.session.commit()
  151. return {
  152. "credential_uid": credential_uid,
  153. "credential": token,
  154. "expires_in": ttl,
  155. "returned_once": True,
  156. }
  157. def authenticate(
  158. self,
  159. token,
  160. *,
  161. connector_id,
  162. connector_version,
  163. source_uid,
  164. business_domain_uid,
  165. environment,
  166. operation,
  167. scope,
  168. ):
  169. scope = validate_machine_scope(scope)
  170. digest = _token_hash(token)
  171. expired = (
  172. self.session.execute(
  173. text("""
  174. SELECT c.uid::text, c.principal_uid::text
  175. FROM public.connector_machine_credentials c
  176. WHERE c.token_hash=:hash AND c.status='active' AND c.expires_at<=CURRENT_TIMESTAMP
  177. FOR UPDATE
  178. """),
  179. {"hash": digest},
  180. )
  181. .mappings()
  182. .one_or_none()
  183. )
  184. if expired is not None:
  185. self.session.execute(
  186. text(
  187. "UPDATE public.connector_machine_credentials SET status='expired',revoked_at=CURRENT_TIMESTAMP WHERE uid=CAST(:uid AS uuid)"
  188. ),
  189. {"uid": expired["uid"]},
  190. )
  191. self._audit(
  192. expired["principal_uid"],
  193. expired["uid"],
  194. "credential_expired_rejected",
  195. None,
  196. False,
  197. )
  198. self.session.commit()
  199. raise ConnectorAuthenticationError(
  200. "machine credential is invalid or expired"
  201. )
  202. row = (
  203. self.session.execute(
  204. text("""
  205. SELECT c.uid::text, c.principal_uid::text, c.use_count, p.connector_id,
  206. p.connector_version, p.source_uid::text, p.business_domain_uid::text,
  207. p.environment, p.allowed_operations, p.allowed_scopes,
  208. p.source_binding_uid::text,p.source_binding_version,
  209. b.approved_base_url,b.allowed_host,b.credential_ref,b.approved_config
  210. FROM public.connector_machine_credentials c
  211. JOIN public.connector_principals p ON c.principal_uid=p.uid
  212. JOIN public.connector_source_bindings b
  213. ON b.uid=p.source_binding_uid
  214. AND b.binding_version=p.source_binding_version
  215. AND b.status='approved'
  216. AND b.connector_id=p.connector_id
  217. AND b.connector_version=p.connector_version
  218. AND b.source_uid=p.source_uid
  219. AND b.business_domain_uid=p.business_domain_uid
  220. AND b.environment=p.environment
  221. WHERE c.token_hash=:hash AND c.status='active'
  222. AND c.expires_at>CURRENT_TIMESTAMP AND p.status='active'
  223. FOR UPDATE OF c
  224. """),
  225. {"hash": digest},
  226. )
  227. .mappings()
  228. .one_or_none()
  229. )
  230. if row is None:
  231. self.session.rollback()
  232. raise ConnectorAuthenticationError(
  233. "machine credential is invalid or expired"
  234. )
  235. if int(row["use_count"]) > 0:
  236. self.session.execute(
  237. text(
  238. "UPDATE public.connector_machine_credentials SET status='replayed', revoked_at=CURRENT_TIMESTAMP WHERE uid=CAST(:uid AS uuid)"
  239. ),
  240. {"uid": row["uid"]},
  241. )
  242. self._audit(
  243. row["principal_uid"],
  244. row["uid"],
  245. "credential_replay_rejected",
  246. None,
  247. False,
  248. )
  249. self.session.commit()
  250. raise ConnectorAuthenticationError("machine credential replay was rejected")
  251. if (
  252. not hmac.compare_digest(row["connector_id"], connector_id)
  253. or not hmac.compare_digest(row["connector_version"], connector_version)
  254. or row["source_uid"] != source_uid
  255. or row["business_domain_uid"] != business_domain_uid
  256. or row["environment"] != environment
  257. or operation not in row["allowed_operations"]
  258. ):
  259. self._audit(
  260. row["principal_uid"],
  261. row["uid"],
  262. "credential_scope_rejected",
  263. None,
  264. False,
  265. )
  266. self.session.commit()
  267. raise ConnectorPermissionError("machine credential scope was rejected")
  268. allowed_scope = validate_machine_scope(row["allowed_scopes"] or {})
  269. if any(
  270. set(scope.get(key, ())) - set(allowed_scope.get(key, ())) for key in scope
  271. ):
  272. self._audit(
  273. row["principal_uid"],
  274. row["uid"],
  275. "credential_scope_rejected",
  276. None,
  277. False,
  278. )
  279. self.session.commit()
  280. raise ConnectorPermissionError(
  281. "machine credential resource scope was rejected"
  282. )
  283. identity = dict(row)
  284. config = dict(identity.pop("approved_config", {}) or {})
  285. if identity["connector_id"] == "rest-catalog":
  286. if not identity.get("source_binding_uid"):
  287. raise ConnectorPermissionError(
  288. "machine credential has no approved source binding"
  289. )
  290. config = {
  291. "base_url": identity.pop("approved_base_url"),
  292. "allowed_host": identity.pop("allowed_host"),
  293. "credential_ref": identity.pop("credential_ref"),
  294. }
  295. else:
  296. identity.pop("approved_base_url", None)
  297. identity.pop("allowed_host", None)
  298. identity.pop("credential_ref", None)
  299. identity["approved_config"] = config
  300. self.session.execute(
  301. text("""
  302. UPDATE public.connector_machine_credentials
  303. SET first_used_at=CURRENT_TIMESTAMP,use_count=use_count+1
  304. WHERE uid=CAST(:uid AS uuid) AND use_count=0 AND status='active'
  305. """),
  306. {"uid": row["uid"]},
  307. )
  308. self._audit(
  309. row["principal_uid"], row["uid"], "credential_authenticated", None, True
  310. )
  311. self.session.commit()
  312. return identity
  313. def revoke(self, credential_uid, actor_uid):
  314. principal_uid = self.session.execute(
  315. text(
  316. "UPDATE public.connector_machine_credentials SET status='revoked', revoked_at=CURRENT_TIMESTAMP WHERE uid=CAST(:uid AS uuid) AND status='active' RETURNING principal_uid::text"
  317. ),
  318. {"uid": credential_uid},
  319. ).scalar_one_or_none()
  320. if principal_uid:
  321. self._audit(
  322. principal_uid, credential_uid, "credential_revoked", actor_uid, True
  323. )
  324. self.session.commit()
  325. return principal_uid is not None
  326. def rotate(self, credential_uid, *, ttl_seconds, actor_uid):
  327. row = self.session.execute(
  328. text(
  329. "UPDATE public.connector_machine_credentials SET status='rotated', revoked_at=CURRENT_TIMESTAMP WHERE uid=CAST(:uid AS uuid) AND status='active' RETURNING principal_uid::text"
  330. ),
  331. {"uid": credential_uid},
  332. ).scalar_one_or_none()
  333. if row is None:
  334. raise ConnectorAuthenticationError("credential cannot be rotated")
  335. self._audit(row, credential_uid, "credential_rotated", actor_uid, True)
  336. return self.issue(
  337. row,
  338. ttl_seconds=ttl_seconds,
  339. actor_uid=actor_uid,
  340. rotated_from=credential_uid,
  341. )
  342. def _audit(self, principal_uid, credential_uid, event_type, actor_uid, success):
  343. self.session.execute(
  344. text("""
  345. INSERT INTO public.connector_audit_events
  346. (uid, principal_uid, credential_uid, event_type, actor_uid, success, safe_detail)
  347. VALUES (CAST(:uid AS uuid), CAST(:principal AS uuid), CAST(:credential AS uuid), :event,
  348. CAST(:actor AS uuid), :success, :detail)
  349. """),
  350. {
  351. "uid": new_governance_uid(),
  352. "principal": principal_uid,
  353. "credential": credential_uid,
  354. "event": event_type,
  355. "actor": actor_uid,
  356. "success": success,
  357. "detail": event_type.replace("_", " "),
  358. },
  359. )
  360. __all__ = ["ConnectorIdentityRepository", "ALLOWED_OPERATIONS"]