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