| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145 |
- """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;
- """)
|