| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379 |
- """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"]
|