"""Harden connector identity bindings and persistent runtime state.""" from alembic import op revision = "20260802_473" down_revision = "20260802_472" branch_labels = None depends_on = None def upgrade() -> None: op.execute(""" DO $$ BEGIN IF EXISTS ( SELECT 1 FROM public.connector_principals p JOIN ( SELECT connector_id FROM public.connector_manifests WHERE status='active' GROUP BY connector_id HAVING COUNT(DISTINCT connector_version)>1 ) ambiguous USING(connector_id) ) THEN RAISE EXCEPTION 'connector 473 upgrade rejected before backfill: multiple active versions require an explicit one-active-version mapping'; END IF; END $$; ALTER TABLE public.connector_principals ADD COLUMN connector_version VARCHAR(40); UPDATE public.connector_principals p SET connector_version=( SELECT m.connector_version FROM public.connector_manifests m WHERE m.connector_id=p.connector_id AND m.status='active' ORDER BY m.created_at DESC LIMIT 1 ); DO $$ BEGIN IF EXISTS (SELECT 1 FROM public.connector_principals WHERE connector_version IS NULL) THEN RAISE EXCEPTION 'connector principal version cannot be backfilled safely'; END IF; END $$; ALTER TABLE public.connector_principals ALTER COLUMN connector_version SET NOT NULL; ALTER TABLE public.connector_principals ADD CONSTRAINT fk_connector_principal_manifest FOREIGN KEY(connector_id,connector_version) REFERENCES public.connector_manifests(connector_id,connector_version); ALTER TABLE public.connector_principals DROP CONSTRAINT uq_connector_principal_binding; ALTER TABLE public.connector_principals ADD CONSTRAINT uq_connector_principal_binding_v2 UNIQUE(connector_id,connector_version,source_uid,business_domain_uid,environment); ALTER TABLE public.connector_runs ADD COLUMN principal_uid UUID REFERENCES public.connector_principals(uid), ADD COLUMN business_domain_uid UUID, ADD COLUMN environment VARCHAR(20), ADD COLUMN process_key VARCHAR(300), ADD COLUMN safe_config JSONB NOT NULL DEFAULT '{}'::jsonb, ADD COLUMN scope JSONB NOT NULL DEFAULT '{}'::jsonb; ALTER TABLE public.connector_runs ADD CONSTRAINT ck_connector_run_environment CHECK(environment IS NULL OR environment IN ('development','staging','production')); ALTER TABLE public.connector_runs ADD CONSTRAINT ck_connector_machine_run_binding CHECK(dry_run OR (principal_uid IS NOT NULL AND business_domain_uid IS NOT NULL AND environment IS NOT NULL AND process_key IS NOT NULL)) NOT VALID; ALTER TABLE public.connector_runs ADD CONSTRAINT ck_connector_run_config_object CHECK(jsonb_typeof(safe_config)='object'); ALTER TABLE public.connector_runs ADD CONSTRAINT ck_connector_run_scope_object CHECK(jsonb_typeof(scope)='object'); CREATE TABLE public.connector_rate_limits( limit_key VARCHAR(300) NOT NULL, window_started_at TIMESTAMPTZ NOT NULL, request_count INTEGER NOT NULL, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY(limit_key,window_started_at), CONSTRAINT ck_connector_rate_count CHECK(request_count BETWEEN 1 AND 30) ); CREATE INDEX ix_connector_rate_limits_updated ON public.connector_rate_limits(updated_at); """) def downgrade() -> None: op.execute(""" DROP TABLE IF EXISTS public.connector_rate_limits; ALTER TABLE public.connector_runs DROP CONSTRAINT IF EXISTS ck_connector_run_scope_object; ALTER TABLE public.connector_runs DROP CONSTRAINT IF EXISTS ck_connector_run_config_object; ALTER TABLE public.connector_runs DROP CONSTRAINT IF EXISTS ck_connector_machine_run_binding; ALTER TABLE public.connector_runs DROP CONSTRAINT IF EXISTS ck_connector_run_environment; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS scope; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS safe_config; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS process_key; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS environment; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS business_domain_uid; ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS principal_uid; ALTER TABLE public.connector_principals DROP CONSTRAINT IF EXISTS uq_connector_principal_binding_v2; ALTER TABLE public.connector_principals DROP CONSTRAINT IF EXISTS fk_connector_principal_manifest; ALTER TABLE public.connector_principals DROP COLUMN IF EXISTS connector_version; ALTER TABLE public.connector_principals ADD CONSTRAINT uq_connector_principal_binding UNIQUE(connector_id,source_uid,business_domain_uid,environment); """)