20260802_475_connector_security_bindings.py 6.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145
  1. """Bind connector destinations, idempotency requests and attempt leases."""
  2. from alembic import op
  3. revision = "20260802_475"
  4. down_revision = "20260802_474"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.execute("""
  9. CREATE TABLE public.connector_source_bindings (
  10. uid UUID NOT NULL,
  11. binding_version INTEGER NOT NULL,
  12. connector_id VARCHAR(64) NOT NULL,
  13. connector_version VARCHAR(40) NOT NULL,
  14. source_uid UUID NOT NULL,
  15. business_domain_uid UUID NOT NULL,
  16. environment VARCHAR(20) NOT NULL,
  17. approved_base_url VARCHAR(2048),
  18. allowed_host VARCHAR(253),
  19. credential_ref VARCHAR(300),
  20. approved_config JSONB NOT NULL DEFAULT '{}'::jsonb,
  21. status VARCHAR(20) NOT NULL,
  22. approved_by UUID NOT NULL,
  23. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  24. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  25. PRIMARY KEY(uid,binding_version),
  26. CONSTRAINT fk_connector_source_binding_manifest
  27. FOREIGN KEY(connector_id,connector_version)
  28. REFERENCES public.connector_manifests(connector_id,connector_version),
  29. CONSTRAINT ck_connector_source_binding_environment
  30. CHECK(environment IN ('development','staging','production')),
  31. CONSTRAINT ck_connector_source_binding_status
  32. CHECK(status IN ('approved','revoked')),
  33. CONSTRAINT ck_connector_source_binding_config
  34. CHECK(jsonb_typeof(approved_config)='object'),
  35. CONSTRAINT ck_connector_source_binding_rest_destination CHECK(
  36. connector_id<>'rest-catalog' OR
  37. (approved_base_url ~ '^https://[^/?#]+(?:/[^?#]*)?$'
  38. AND allowed_host IS NOT NULL AND credential_ref IS NOT NULL)
  39. )
  40. );
  41. CREATE UNIQUE INDEX uq_connector_source_binding_active
  42. ON public.connector_source_bindings(
  43. connector_id,connector_version,source_uid,business_domain_uid,environment
  44. ) WHERE status='approved';
  45. CREATE INDEX ix_connector_source_binding_lookup
  46. ON public.connector_source_bindings(source_uid,business_domain_uid,environment,status);
  47. ALTER TABLE public.connector_principals
  48. ADD COLUMN source_binding_uid UUID,
  49. ADD COLUMN source_binding_version INTEGER;
  50. ALTER TABLE public.connector_principals
  51. ADD CONSTRAINT fk_connector_principal_source_binding
  52. FOREIGN KEY(source_binding_uid,source_binding_version)
  53. REFERENCES public.connector_source_bindings(uid,binding_version);
  54. ALTER TABLE public.connector_runs
  55. ADD COLUMN request_hash CHAR(64),
  56. ADD COLUMN client_hint_hash CHAR(64),
  57. ADD COLUMN source_binding_uid UUID,
  58. ADD COLUMN source_binding_version INTEGER,
  59. ADD COLUMN attempt_lease_token UUID,
  60. ADD COLUMN cancel_requested BOOLEAN NOT NULL DEFAULT FALSE;
  61. ALTER TABLE public.connector_runs
  62. ADD CONSTRAINT ck_connector_run_request_hash
  63. CHECK(request_hash IS NOT NULL) NOT VALID;
  64. ALTER TABLE public.connector_runs
  65. ADD CONSTRAINT fk_connector_run_source_binding
  66. FOREIGN KEY(source_binding_uid,source_binding_version)
  67. REFERENCES public.connector_source_bindings(uid,binding_version);
  68. CREATE INDEX ix_connector_runs_request_hash ON public.connector_runs(request_hash);
  69. CREATE UNIQUE INDEX uq_connector_runs_client_hint
  70. ON public.connector_runs(client_hint_hash) WHERE client_hint_hash IS NOT NULL;
  71. CREATE INDEX ix_connector_runs_binding
  72. ON public.connector_runs(source_binding_uid,source_binding_version,created_at DESC);
  73. ALTER TABLE public.connector_run_attempts ADD COLUMN lease_token UUID;
  74. CREATE UNIQUE INDEX uq_connector_attempt_lease
  75. ON public.connector_run_attempts(run_uid,lease_token) WHERE lease_token IS NOT NULL;
  76. CREATE OR REPLACE FUNCTION public.validate_connector_binding()
  77. RETURNS trigger LANGUAGE plpgsql AS $$
  78. BEGIN
  79. IF NEW.source_binding_uid IS NULL OR NEW.source_binding_version IS NULL THEN
  80. IF NEW.connector_id='rest-catalog' THEN
  81. RAISE EXCEPTION 'REST connector principal requires an approved source binding';
  82. END IF;
  83. RETURN NEW;
  84. END IF;
  85. IF NOT EXISTS (
  86. SELECT 1 FROM public.connector_source_bindings b
  87. WHERE b.uid=NEW.source_binding_uid
  88. AND b.binding_version=NEW.source_binding_version
  89. AND b.status='approved'
  90. AND b.connector_id=NEW.connector_id
  91. AND b.connector_version=NEW.connector_version
  92. AND b.source_uid=NEW.source_uid
  93. AND b.business_domain_uid=NEW.business_domain_uid
  94. AND b.environment=NEW.environment
  95. ) THEN
  96. RAISE EXCEPTION 'connector principal source binding mismatch';
  97. END IF;
  98. RETURN NEW;
  99. END $$;
  100. CREATE TRIGGER trg_connector_principal_binding
  101. BEFORE INSERT OR UPDATE OF connector_id,connector_version,source_uid,
  102. business_domain_uid,environment,source_binding_uid,source_binding_version
  103. ON public.connector_principals
  104. FOR EACH ROW EXECUTE FUNCTION public.validate_connector_binding();
  105. """)
  106. def downgrade() -> None:
  107. op.execute("""
  108. DO $$ BEGIN
  109. IF EXISTS (SELECT 1 FROM public.connector_source_bindings)
  110. OR EXISTS (SELECT 1 FROM public.connector_principals WHERE source_binding_uid IS NOT NULL)
  111. OR EXISTS (SELECT 1 FROM public.connector_runs WHERE source_binding_uid IS NOT NULL) THEN
  112. RAISE EXCEPTION
  113. 'connector downgrade below 475 rejected: revoke and explicitly remove source bindings first';
  114. END IF;
  115. END $$;
  116. DROP TRIGGER IF EXISTS trg_connector_principal_binding ON public.connector_principals;
  117. DROP FUNCTION IF EXISTS public.validate_connector_binding();
  118. DROP INDEX IF EXISTS public.uq_connector_attempt_lease;
  119. ALTER TABLE public.connector_run_attempts DROP COLUMN IF EXISTS lease_token;
  120. DROP INDEX IF EXISTS public.ix_connector_runs_binding;
  121. DROP INDEX IF EXISTS public.ix_connector_runs_request_hash;
  122. DROP INDEX IF EXISTS public.uq_connector_runs_client_hint;
  123. ALTER TABLE public.connector_runs DROP CONSTRAINT IF EXISTS fk_connector_run_source_binding;
  124. ALTER TABLE public.connector_runs DROP CONSTRAINT IF EXISTS ck_connector_run_request_hash;
  125. ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS cancel_requested;
  126. ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS attempt_lease_token;
  127. ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS source_binding_version;
  128. ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS source_binding_uid;
  129. ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS request_hash;
  130. ALTER TABLE public.connector_runs DROP COLUMN IF EXISTS client_hint_hash;
  131. ALTER TABLE public.connector_principals DROP CONSTRAINT IF EXISTS fk_connector_principal_source_binding;
  132. ALTER TABLE public.connector_principals DROP COLUMN IF EXISTS source_binding_version;
  133. ALTER TABLE public.connector_principals DROP COLUMN IF EXISTS source_binding_uid;
  134. DROP TABLE public.connector_source_bindings;
  135. """)