"""Bind connector destinations, idempotency requests and attempt leases.""" from alembic import op revision = "20260802_475" down_revision = "20260802_474" branch_labels = None depends_on = None def upgrade() -> None: op.execute(""" CREATE TABLE public.connector_source_bindings ( uid UUID NOT NULL, binding_version INTEGER NOT NULL, connector_id VARCHAR(64) NOT NULL, connector_version VARCHAR(40) NOT NULL, source_uid UUID NOT NULL, business_domain_uid UUID NOT NULL, environment VARCHAR(20) NOT NULL, approved_base_url VARCHAR(2048), allowed_host VARCHAR(253), credential_ref VARCHAR(300), approved_config JSONB NOT NULL DEFAULT '{}'::jsonb, status VARCHAR(20) NOT NULL, approved_by UUID NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY(uid,binding_version), CONSTRAINT fk_connector_source_binding_manifest FOREIGN KEY(connector_id,connector_version) REFERENCES public.connector_manifests(connector_id,connector_version), CONSTRAINT ck_connector_source_binding_environment CHECK(environment IN ('development','staging','production')), CONSTRAINT ck_connector_source_binding_status CHECK(status IN ('approved','revoked')), CONSTRAINT ck_connector_source_binding_config CHECK(jsonb_typeof(approved_config)='object'), CONSTRAINT ck_connector_source_binding_rest_destination CHECK( connector_id<>'rest-catalog' OR (approved_base_url ~ '^https://[^/?#]+(?:/[^?#]*)?$' AND allowed_host IS NOT NULL AND credential_ref IS NOT NULL) ) ); CREATE UNIQUE INDEX uq_connector_source_binding_active ON public.connector_source_bindings( connector_id,connector_version,source_uid,business_domain_uid,environment ) WHERE status='approved'; CREATE INDEX ix_connector_source_binding_lookup ON public.connector_source_bindings(source_uid,business_domain_uid,environment,status); ALTER TABLE public.connector_principals ADD COLUMN source_binding_uid UUID, ADD COLUMN source_binding_version INTEGER; ALTER TABLE public.connector_principals ADD CONSTRAINT fk_connector_principal_source_binding FOREIGN KEY(source_binding_uid,source_binding_version) REFERENCES public.connector_source_bindings(uid,binding_version); ALTER TABLE public.connector_runs ADD COLUMN request_hash CHAR(64), ADD COLUMN client_hint_hash CHAR(64), ADD COLUMN source_binding_uid UUID, ADD COLUMN source_binding_version INTEGER, ADD COLUMN attempt_lease_token UUID, ADD COLUMN cancel_requested BOOLEAN NOT NULL DEFAULT FALSE; ALTER TABLE public.connector_runs ADD CONSTRAINT ck_connector_run_request_hash CHECK(request_hash IS NOT NULL) NOT VALID; ALTER TABLE public.connector_runs ADD CONSTRAINT fk_connector_run_source_binding FOREIGN KEY(source_binding_uid,source_binding_version) REFERENCES public.connector_source_bindings(uid,binding_version); CREATE INDEX ix_connector_runs_request_hash ON public.connector_runs(request_hash); CREATE UNIQUE INDEX uq_connector_runs_client_hint ON public.connector_runs(client_hint_hash) WHERE client_hint_hash IS NOT NULL; CREATE INDEX ix_connector_runs_binding ON public.connector_runs(source_binding_uid,source_binding_version,created_at DESC); ALTER TABLE public.connector_run_attempts ADD COLUMN lease_token UUID; CREATE UNIQUE INDEX uq_connector_attempt_lease ON public.connector_run_attempts(run_uid,lease_token) WHERE lease_token IS NOT NULL; CREATE OR REPLACE FUNCTION public.validate_connector_binding() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN IF NEW.source_binding_uid IS NULL OR NEW.source_binding_version IS NULL THEN IF NEW.connector_id='rest-catalog' THEN RAISE EXCEPTION 'REST connector principal requires an approved source binding'; END IF; RETURN NEW; END IF; IF NOT EXISTS ( SELECT 1 FROM public.connector_source_bindings b WHERE b.uid=NEW.source_binding_uid AND b.binding_version=NEW.source_binding_version AND b.status='approved' AND b.connector_id=NEW.connector_id AND b.connector_version=NEW.connector_version AND b.source_uid=NEW.source_uid AND b.business_domain_uid=NEW.business_domain_uid AND b.environment=NEW.environment ) THEN RAISE EXCEPTION 'connector principal source binding mismatch'; END IF; RETURN NEW; END $$; CREATE TRIGGER trg_connector_principal_binding BEFORE INSERT OR UPDATE OF connector_id,connector_version,source_uid, business_domain_uid,environment,source_binding_uid,source_binding_version ON public.connector_principals FOR EACH ROW EXECUTE FUNCTION public.validate_connector_binding(); """) def downgrade() -> None: op.execute(""" DO $$ BEGIN IF EXISTS (SELECT 1 FROM public.connector_source_bindings) OR EXISTS (SELECT 1 FROM public.connector_principals WHERE source_binding_uid IS NOT NULL) OR EXISTS (SELECT 1 FROM public.connector_runs WHERE source_binding_uid IS NOT NULL) THEN RAISE EXCEPTION 'connector downgrade below 475 rejected: revoke and explicitly remove source bindings first'; END IF; END $$; DROP TRIGGER IF EXISTS trg_connector_principal_binding ON public.connector_principals; DROP FUNCTION IF EXISTS public.validate_connector_binding(); DROP INDEX IF EXISTS public.uq_connector_attempt_lease; ALTER TABLE public.connector_run_attempts DROP COLUMN IF EXISTS lease_token; DROP INDEX IF EXISTS public.ix_connector_runs_binding; DROP INDEX IF EXISTS public.ix_connector_runs_request_hash; DROP INDEX IF EXISTS public.uq_connector_runs_client_hint; ALTER TABLE public.connector_runs DROP CONSTRAINT IF EXISTS fk_connector_run_source_binding; ALTER TABLE public.connector_runs DROP CONSTRAINT IF EXISTS ck_connector_run_request_hash; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS cancel_requested; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS attempt_lease_token; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS source_binding_version; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS source_binding_uid; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS request_hash; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS client_hint_hash; ALTER TABLE public.connector_principals DROP CONSTRAINT IF EXISTS fk_connector_principal_source_binding; ALTER TABLE public.connector_principals DROP COLUMN IF EXISTS source_binding_version; ALTER TABLE public.connector_principals DROP COLUMN IF EXISTS source_binding_uid; DROP TABLE public.connector_source_bindings; """)