"""Server-approved connector source and destination bindings.""" from __future__ import annotations import json from urllib.parse import urlsplit from sqlalchemy import text from app.core.common.identifiers import new_governance_uid from app.core.connectors.errors import ConnectorConfigurationError from app.core.connectors.sdk import SECRET_REF_PATTERN, _walk_secrets ENVIRONMENTS = frozenset({"development", "staging", "production"}) def validate_approved_config(connector_id, config): if not isinstance(config, dict): raise ConnectorConfigurationError("approved connector config must be an object") _walk_secrets(config, "approved_config") if connector_id != "rest-catalog": return dict(config) if set(config) != {"base_url", "allowed_host", "credential_ref"}: raise ConnectorConfigurationError("REST approved config is incomplete") parsed = urlsplit(str(config["base_url"])) host = str(config["allowed_host"]).lower().rstrip(".") if ( parsed.scheme != "https" or not parsed.hostname or parsed.username or parsed.password or parsed.query or parsed.fragment or parsed.hostname.lower().rstrip(".") != host ): raise ConnectorConfigurationError("REST approved destination is invalid") credential_ref = str(config["credential_ref"]) if not SECRET_REF_PATTERN.fullmatch(credential_ref): raise ConnectorConfigurationError("REST approved credential reference is invalid") return { "base_url": str(config["base_url"]).rstrip("/"), "allowed_host": host, "credential_ref": credential_ref, } class ConnectorSourceBindingRepository: def __init__(self, session): self.session = session def approve( self, *, connector_id, connector_version, source_uid, business_domain_uid, environment, approved_config, approved_by, binding_uid=None, ): if environment not in ENVIRONMENTS: raise ConnectorConfigurationError("connector binding environment is invalid") config = validate_approved_config(connector_id, approved_config) uid = binding_uid or new_governance_uid() if binding_uid: self.session.execute( text("SELECT pg_advisory_xact_lock(hashtext(:lock_key))"), {"lock_key": f"connector-binding:{uid}"}, ) current = self.session.execute( text(""" SELECT COALESCE(MAX(binding_version),0) FROM public.connector_source_bindings WHERE uid=CAST(:uid AS uuid) """), {"uid": uid}, ).scalar_one() self.session.execute( text(""" UPDATE public.connector_source_bindings SET status='revoked',updated_at=CURRENT_TIMESTAMP WHERE uid=CAST(:uid AS uuid) AND status='approved' """), {"uid": uid}, ) version = int(current) + 1 else: version = 1 rest = config if connector_id == "rest-catalog" else {} self.session.execute( text(""" INSERT INTO public.connector_source_bindings (uid,binding_version,connector_id,connector_version,source_uid, business_domain_uid,environment,approved_base_url,allowed_host, credential_ref,approved_config,status,approved_by) VALUES(CAST(:uid AS uuid),:binding_version,:connector,:version, CAST(:source AS uuid),CAST(:domain AS uuid),:environment, :base_url,:allowed_host,:credential_ref,CAST(:config AS jsonb), 'approved',CAST(:approved_by AS uuid)) """), { "uid": uid, "binding_version": version, "connector": connector_id, "version": connector_version, "source": source_uid, "domain": business_domain_uid, "environment": environment, "base_url": rest.get("base_url"), "allowed_host": rest.get("allowed_host"), "credential_ref": rest.get("credential_ref"), "config": json.dumps(config, sort_keys=True), "approved_by": approved_by, }, ) rebound_principals = 0 if binding_uid and int(current) > 0: rebound_principals = self.session.execute( text(""" UPDATE public.connector_principals SET source_binding_version=:binding_version WHERE source_binding_uid=CAST(:uid AS uuid) AND source_binding_version=:old_version """), { "uid": uid, "old_version": int(current), "binding_version": version, }, ).rowcount self.session.commit() return { "binding_uid": uid, "binding_version": version, "status": "approved", "rebound_principals": rebound_principals, } def revoke(self, binding_uid, approved_by): principals = self.session.execute( text(""" UPDATE public.connector_principals SET status='revoked' WHERE source_binding_uid=CAST(:uid AS uuid) AND status='active' RETURNING uid """), {"uid": binding_uid}, ).all() if principals: principal_uids = [str(row[0]) for row in principals] self.session.execute( text(""" UPDATE public.connector_machine_credentials SET status='revoked',revoked_at=CURRENT_TIMESTAMP WHERE principal_uid=ANY(CAST(:principals AS uuid[])) AND status='active' """), {"principals": principal_uids}, ) changed = self.session.execute( text(""" UPDATE public.connector_source_bindings SET status='revoked',approved_by=CAST(:actor AS uuid),updated_at=CURRENT_TIMESTAMP WHERE uid=CAST(:uid AS uuid) AND status='approved' """), {"uid": binding_uid, "actor": approved_by}, ).rowcount self.session.commit() return {"revoked": bool(changed), "principals_deactivated": len(principals)} def get_approved(self, binding_uid, binding_version=None): version_clause = ( "AND binding_version=:binding_version" if binding_version is not None else "" ) row = ( self.session.execute( text(f""" SELECT uid::text,binding_version,connector_id,connector_version, source_uid::text,business_domain_uid::text,environment, approved_base_url,allowed_host,credential_ref,approved_config,status FROM public.connector_source_bindings WHERE uid=CAST(:uid AS uuid) AND status='approved' {version_clause} ORDER BY binding_version DESC LIMIT 1 """), {"uid": binding_uid, "binding_version": binding_version}, ) .mappings() .one_or_none() ) if row is None: raise ConnectorConfigurationError("approved connector binding was not found") result = dict(row) config = dict(result.pop("approved_config") or {}) if result["connector_id"] == "rest-catalog": config = { "base_url": result.pop("approved_base_url"), "allowed_host": result.pop("allowed_host"), "credential_ref": result.pop("credential_ref"), } else: result.pop("approved_base_url", None) result.pop("allowed_host", None) result.pop("credential_ref", None) result["approved_config"] = config return result def list_public(self, limit=100): rows = self.session.execute( text(""" SELECT uid::text,binding_version,connector_id,connector_version, source_uid::text,business_domain_uid::text,environment, approved_base_url,allowed_host,status,approved_by::text, created_at,updated_at FROM public.connector_source_bindings ORDER BY created_at DESC LIMIT :limit """), {"limit": min(max(int(limit), 1), 200)}, ).mappings() return [dict(row) for row in rows] __all__ = ["ConnectorSourceBindingRepository", "validate_approved_config"]