"""Short-lived, one-time-issued connector machine credentials.""" from __future__ import annotations import hashlib import hmac import secrets from collections.abc import Mapping from sqlalchemy import text from app.core.common.identifiers import new_governance_uid from app.core.connectors.errors import ( ConnectorAuthenticationError, ConnectorConfigurationError, ConnectorPermissionError, ) ALLOWED_OPERATIONS = frozenset( { "discover", "snapshot", "incremental", "lineage", "profile", "cancel", "resume", "evidence", } ) ALLOWED_SCOPE_KEYS = frozenset( {"include_schemas", "exclude_schemas", "include_tables", "exclude_tables"} ) def validate_machine_scope(scope): if not isinstance(scope, Mapping) or set(scope) - ALLOWED_SCOPE_KEYS: raise ConnectorConfigurationError("machine identity scope is invalid") normalized = {} for key, values in scope.items(): if not isinstance(values, list) or any( not isinstance(value, str) or not value.strip() for value in values ): raise ConnectorConfigurationError( "machine identity scope values are invalid" ) normalized[key] = tuple(value.strip() for value in values) return normalized def _token_hash(token): return hashlib.sha256(str(token).encode()).hexdigest() class ConnectorIdentityRepository: def __init__(self, session): self.session = session def create_principal( self, *, connector_id, connector_version, source_uid, business_domain_uid, environment, operations, scopes, actor_uid, source_binding_uid=None, source_binding_version=None, ): operations = tuple(sorted(set(operations))) if not operations or set(operations) - ALLOWED_OPERATIONS: raise ConnectorConfigurationError("machine identity operations are invalid") scopes = validate_machine_scope(scopes) if environment not in {"development", "staging", "production"}: raise ConnectorConfigurationError("machine identity environment is invalid") if not source_binding_uid or source_binding_version is None: raise ConnectorConfigurationError( "enterprise connector principal requires an approved source binding" ) uid = new_governance_uid() self.session.execute( text(""" INSERT INTO public.connector_principals (uid, connector_id, connector_version, source_uid, business_domain_uid, environment, allowed_operations, allowed_scopes, status, created_by, source_binding_uid,source_binding_version) VALUES (CAST(:uid AS uuid), :connector, :version, CAST(:source AS uuid), CAST(:domain AS uuid), :environment, CAST(:operations AS text[]), CAST(:scopes AS jsonb), 'active', CAST(:actor AS uuid), CAST(:binding AS uuid),:binding_version) """), { "uid": uid, "connector": connector_id, "version": connector_version, "source": source_uid, "domain": business_domain_uid, "environment": environment, "operations": list(operations), "scopes": __import__("json").dumps(scopes), "actor": actor_uid, "binding": source_binding_uid, "binding_version": source_binding_version, }, ) self._audit(uid, None, "principal_created", actor_uid, True) self.session.commit() return uid def issue(self, principal_uid, *, ttl_seconds, actor_uid, rotated_from=None): ttl = int(ttl_seconds) if ttl < 60 or ttl > 900: raise ConnectorConfigurationError( "credential TTL must be between 60 and 900 seconds" ) self.session.execute( text("SELECT pg_advisory_xact_lock(hashtext(:key))"), {"key": f"connector-principal:{principal_uid}"}, ) principal = self.session.execute( text(""" SELECT p.status FROM public.connector_principals p JOIN public.connector_source_bindings b ON b.uid=p.source_binding_uid AND b.binding_version=p.source_binding_version AND b.status='approved' AND b.connector_id=p.connector_id AND b.connector_version=p.connector_version AND b.source_uid=p.source_uid AND b.business_domain_uid=p.business_domain_uid AND b.environment=p.environment WHERE p.uid=CAST(:uid AS uuid) FOR UPDATE OF p """), {"uid": principal_uid}, ).scalar_one_or_none() if principal != "active": raise ConnectorAuthenticationError( "connector principal or approved binding is not active" ) token = "dopc_" + secrets.token_urlsafe(32) credential_uid = new_governance_uid() self.session.execute( text(""" INSERT INTO public.connector_machine_credentials (uid, principal_uid, token_hash, status, expires_at, rotated_from_uid, issued_by) VALUES (CAST(:uid AS uuid), CAST(:principal AS uuid), :hash, 'active', CURRENT_TIMESTAMP + (:ttl * INTERVAL '1 second'), CAST(:rotated AS uuid), CAST(:actor AS uuid)) """), { "uid": credential_uid, "principal": principal_uid, "hash": _token_hash(token), "ttl": ttl, "rotated": rotated_from, "actor": actor_uid, }, ) self._audit(principal_uid, credential_uid, "credential_issued", actor_uid, True) self.session.commit() return { "credential_uid": credential_uid, "credential": token, "expires_in": ttl, "returned_once": True, } def authenticate( self, token, *, connector_id, connector_version, source_uid, business_domain_uid, environment, operation, scope, ): scope = validate_machine_scope(scope) digest = _token_hash(token) expired = ( self.session.execute( text(""" SELECT c.uid::text, c.principal_uid::text FROM public.connector_machine_credentials c WHERE c.token_hash=:hash AND c.status='active' AND c.expires_at<=CURRENT_TIMESTAMP FOR UPDATE """), {"hash": digest}, ) .mappings() .one_or_none() ) if expired is not None: self.session.execute( text( "UPDATE public.connector_machine_credentials SET status='expired',revoked_at=CURRENT_TIMESTAMP WHERE uid=CAST(:uid AS uuid)" ), {"uid": expired["uid"]}, ) self._audit( expired["principal_uid"], expired["uid"], "credential_expired_rejected", None, False, ) self.session.commit() raise ConnectorAuthenticationError( "machine credential is invalid or expired" ) row = ( self.session.execute( text(""" SELECT c.uid::text, c.principal_uid::text, c.use_count, p.connector_id, p.connector_version, p.source_uid::text, p.business_domain_uid::text, p.environment, p.allowed_operations, p.allowed_scopes, p.source_binding_uid::text,p.source_binding_version, b.approved_base_url,b.allowed_host,b.credential_ref,b.approved_config FROM public.connector_machine_credentials c JOIN public.connector_principals p ON c.principal_uid=p.uid JOIN public.connector_source_bindings b ON b.uid=p.source_binding_uid AND b.binding_version=p.source_binding_version AND b.status='approved' AND b.connector_id=p.connector_id AND b.connector_version=p.connector_version AND b.source_uid=p.source_uid AND b.business_domain_uid=p.business_domain_uid AND b.environment=p.environment WHERE c.token_hash=:hash AND c.status='active' AND c.expires_at>CURRENT_TIMESTAMP AND p.status='active' FOR UPDATE OF c """), {"hash": digest}, ) .mappings() .one_or_none() ) if row is None: self.session.rollback() raise ConnectorAuthenticationError( "machine credential is invalid or expired" ) if int(row["use_count"]) > 0: self.session.execute( text( "UPDATE public.connector_machine_credentials SET status='replayed', revoked_at=CURRENT_TIMESTAMP WHERE uid=CAST(:uid AS uuid)" ), {"uid": row["uid"]}, ) self._audit( row["principal_uid"], row["uid"], "credential_replay_rejected", None, False, ) self.session.commit() raise ConnectorAuthenticationError("machine credential replay was rejected") if ( not hmac.compare_digest(row["connector_id"], connector_id) or not hmac.compare_digest(row["connector_version"], connector_version) or row["source_uid"] != source_uid or row["business_domain_uid"] != business_domain_uid or row["environment"] != environment or operation not in row["allowed_operations"] ): self._audit( row["principal_uid"], row["uid"], "credential_scope_rejected", None, False, ) self.session.commit() raise ConnectorPermissionError("machine credential scope was rejected") allowed_scope = validate_machine_scope(row["allowed_scopes"] or {}) if any( set(scope.get(key, ())) - set(allowed_scope.get(key, ())) for key in scope ): self._audit( row["principal_uid"], row["uid"], "credential_scope_rejected", None, False, ) self.session.commit() raise ConnectorPermissionError( "machine credential resource scope was rejected" ) identity = dict(row) config = dict(identity.pop("approved_config", {}) or {}) if identity["connector_id"] == "rest-catalog": if not identity.get("source_binding_uid"): raise ConnectorPermissionError( "machine credential has no approved source binding" ) config = { "base_url": identity.pop("approved_base_url"), "allowed_host": identity.pop("allowed_host"), "credential_ref": identity.pop("credential_ref"), } else: identity.pop("approved_base_url", None) identity.pop("allowed_host", None) identity.pop("credential_ref", None) identity["approved_config"] = config self.session.execute( text(""" UPDATE public.connector_machine_credentials SET first_used_at=CURRENT_TIMESTAMP,use_count=use_count+1 WHERE uid=CAST(:uid AS uuid) AND use_count=0 AND status='active' """), {"uid": row["uid"]}, ) self._audit( row["principal_uid"], row["uid"], "credential_authenticated", None, True ) self.session.commit() return identity def revoke(self, credential_uid, actor_uid): principal_uid = self.session.execute( text( "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" ), {"uid": credential_uid}, ).scalar_one_or_none() if principal_uid: self._audit( principal_uid, credential_uid, "credential_revoked", actor_uid, True ) self.session.commit() return principal_uid is not None def rotate(self, credential_uid, *, ttl_seconds, actor_uid): row = self.session.execute( text( "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" ), {"uid": credential_uid}, ).scalar_one_or_none() if row is None: raise ConnectorAuthenticationError("credential cannot be rotated") self._audit(row, credential_uid, "credential_rotated", actor_uid, True) return self.issue( row, ttl_seconds=ttl_seconds, actor_uid=actor_uid, rotated_from=credential_uid, ) def _audit(self, principal_uid, credential_uid, event_type, actor_uid, success): self.session.execute( text(""" INSERT INTO public.connector_audit_events (uid, principal_uid, credential_uid, event_type, actor_uid, success, safe_detail) VALUES (CAST(:uid AS uuid), CAST(:principal AS uuid), CAST(:credential AS uuid), :event, CAST(:actor AS uuid), :success, :detail) """), { "uid": new_governance_uid(), "principal": principal_uid, "credential": credential_uid, "event": event_type, "actor": actor_uid, "success": success, "detail": event_type.replace("_", " "), }, ) __all__ = ["ConnectorIdentityRepository", "ALLOWED_OPERATIONS"]