20260724_240_data_factory_deployments.py 9.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214
  1. """Add recoverable immutable Data Factory deployment lifecycle."""
  2. from alembic import op
  3. revision = "20260724_240"
  4. down_revision = "20260724_230"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.execute(
  9. """
  10. ALTER TABLE public.dataflow_deployments
  11. ADD COLUMN dataflow_uid UUID,
  12. ADD COLUMN package JSONB,
  13. ADD COLUMN package_hash CHAR(64),
  14. ADD COLUMN standard_version_ids JSONB NOT NULL DEFAULT '[]'::jsonb,
  15. ADD COLUMN rule_version_ids JSONB NOT NULL DEFAULT '[]'::jsonb,
  16. ADD COLUMN binding_snapshot JSONB,
  17. ADD COLUMN binding_hash CHAR(64),
  18. ADD COLUMN schema_snapshots JSONB,
  19. ADD COLUMN schema_snapshot_hash CHAR(64),
  20. ADD COLUMN physical_plan_hashes JSONB NOT NULL DEFAULT '[]'::jsonb,
  21. ADD COLUMN workflow_spec JSONB,
  22. ADD COLUMN workflow_spec_hash CHAR(64),
  23. ADD COLUMN schedule_snapshot JSONB,
  24. ADD COLUMN schedule_hash CHAR(64),
  25. ADD COLUMN engine_namespace VARCHAR(200),
  26. ADD COLUMN engine_definition_id VARCHAR(200),
  27. ADD COLUMN engine_revision VARCHAR(200),
  28. ADD COLUMN engine_definition_hash CHAR(64),
  29. ADD COLUMN deploy_receipt JSONB NOT NULL DEFAULT '{}'::jsonb,
  30. ADD COLUMN canary_evidence_id UUID,
  31. ADD COLUMN previous_active_deployment_id UUID
  32. REFERENCES public.dataflow_deployments(id) ON DELETE RESTRICT,
  33. ADD COLUMN created_by UUID REFERENCES public.users(id) ON DELETE SET NULL,
  34. ADD COLUMN create_reason VARCHAR(1000),
  35. ADD COLUMN create_idempotency_key VARCHAR(200),
  36. ADD COLUMN correlation_id UUID,
  37. ADD COLUMN lock_version INTEGER NOT NULL DEFAULT 0,
  38. ADD COLUMN failure_code VARCHAR(100),
  39. ADD COLUMN superseded_at TIMESTAMPTZ,
  40. ADD COLUMN rolled_back_at TIMESTAMPTZ;
  41. UPDATE public.dataflow_deployments d
  42. SET dataflow_uid = v.dataflow_uid
  43. FROM public.dataflow_versions v
  44. WHERE v.id = d.dataflow_version_id
  45. AND d.dataflow_uid IS NULL;
  46. ALTER TABLE public.dataflow_deployments
  47. ADD CONSTRAINT ck_dataflow_deployment_lock_version
  48. CHECK (lock_version >= 0),
  49. ADD CONSTRAINT ck_dataflow_deployment_hash_shape
  50. CHECK (
  51. (package_hash IS NULL OR package_hash ~ '^[0-9a-f]{64}$')
  52. AND (binding_hash IS NULL OR binding_hash ~ '^[0-9a-f]{64}$')
  53. AND (
  54. schema_snapshot_hash IS NULL
  55. OR schema_snapshot_hash ~ '^[0-9a-f]{64}$'
  56. )
  57. AND (
  58. workflow_spec_hash IS NULL
  59. OR workflow_spec_hash ~ '^[0-9a-f]{64}$'
  60. )
  61. AND (
  62. schedule_hash IS NULL
  63. OR schedule_hash ~ '^[0-9a-f]{64}$'
  64. )
  65. AND (
  66. engine_definition_hash IS NULL
  67. OR engine_definition_hash ~ '^[0-9a-f]{64}$'
  68. )
  69. ),
  70. ADD CONSTRAINT ck_dataflow_deployment_snapshot_complete
  71. CHECK (
  72. package IS NULL
  73. OR (
  74. dataflow_uid IS NOT NULL
  75. AND package_hash IS NOT NULL
  76. AND binding_snapshot IS NOT NULL
  77. AND binding_hash IS NOT NULL
  78. AND schema_snapshots IS NOT NULL
  79. AND schema_snapshot_hash IS NOT NULL
  80. AND jsonb_array_length(physical_plan_hashes) > 0
  81. AND workflow_spec IS NOT NULL
  82. AND workflow_spec_hash IS NOT NULL
  83. AND schedule_snapshot IS NOT NULL
  84. AND schedule_hash IS NOT NULL
  85. AND created_by IS NOT NULL
  86. AND create_reason IS NOT NULL
  87. AND create_idempotency_key IS NOT NULL
  88. AND correlation_id IS NOT NULL
  89. )
  90. );
  91. ALTER TABLE public.dataflow_deployments
  92. DROP CONSTRAINT IF EXISTS
  93. dataflow_deployments_dataflow_version_id_environment_key;
  94. DROP INDEX IF EXISTS public.uq_dataflow_deployment_active_environment;
  95. CREATE UNIQUE INDEX uq_dataflow_deployment_active_environment
  96. ON public.dataflow_deployments(dataflow_uid, environment)
  97. WHERE status = 'active';
  98. CREATE UNIQUE INDEX uq_dataflow_deployment_create_idempotency
  99. ON public.dataflow_deployments(
  100. created_by, environment, create_idempotency_key
  101. )
  102. WHERE create_idempotency_key IS NOT NULL;
  103. CREATE INDEX idx_dataflow_deployment_factory_list
  104. ON public.dataflow_deployments(
  105. environment, updated_at DESC, dataflow_uid
  106. )
  107. WHERE package IS NOT NULL;
  108. CREATE TABLE public.dataflow_canary_evidence (
  109. id UUID PRIMARY KEY,
  110. deployment_id UUID NOT NULL
  111. REFERENCES public.dataflow_deployments(id) ON DELETE RESTRICT,
  112. execution_id VARCHAR(200) NOT NULL,
  113. status VARCHAR(20) NOT NULL CHECK (
  114. status IN ('started','passed','failed','unknown')
  115. ),
  116. package_hash CHAR(64) NOT NULL,
  117. binding_hash CHAR(64) NOT NULL,
  118. schema_snapshot_hash CHAR(64) NOT NULL,
  119. physical_plan_hashes JSONB NOT NULL,
  120. workflow_spec_hash CHAR(64) NOT NULL,
  121. schedule_hash CHAR(64) NOT NULL,
  122. engine_definition_hash CHAR(64) NOT NULL,
  123. verified_by UUID NOT NULL REFERENCES public.users(id) ON DELETE RESTRICT,
  124. engine_response JSONB NOT NULL DEFAULT '{}'::jsonb,
  125. expires_at TIMESTAMPTZ NOT NULL,
  126. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  127. UNIQUE (deployment_id, execution_id)
  128. );
  129. CREATE INDEX idx_dataflow_canary_activation_gate
  130. ON public.dataflow_canary_evidence(
  131. deployment_id, status, expires_at DESC
  132. );
  133. ALTER TABLE public.dataflow_deployments
  134. ADD CONSTRAINT fk_dataflow_deployment_canary_evidence
  135. FOREIGN KEY (canary_evidence_id)
  136. REFERENCES public.dataflow_canary_evidence(id)
  137. ON DELETE RESTRICT;
  138. CREATE TABLE public.dataflow_deployment_operations (
  139. id UUID PRIMARY KEY,
  140. deployment_id UUID NOT NULL
  141. REFERENCES public.dataflow_deployments(id) ON DELETE RESTRICT,
  142. action VARCHAR(30) NOT NULL CHECK (
  143. action IN (
  144. 'deploy_disabled','run_canary','activate','rollback'
  145. )
  146. ),
  147. idempotency_key VARCHAR(200) NOT NULL,
  148. status VARCHAR(20) NOT NULL CHECK (
  149. status IN ('claimed','completed','failed','unknown')
  150. ),
  151. actor_uid UUID NOT NULL REFERENCES public.users(id) ON DELETE RESTRICT,
  152. correlation_id UUID NOT NULL,
  153. reason VARCHAR(1000) NOT NULL,
  154. request JSONB NOT NULL DEFAULT '{}'::jsonb,
  155. request_hash CHAR(64) NOT NULL,
  156. result JSONB,
  157. result_hash CHAR(64),
  158. error_code VARCHAR(100),
  159. started_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  160. completed_at TIMESTAMPTZ,
  161. failed_at TIMESTAMPTZ,
  162. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  163. UNIQUE (deployment_id, action, idempotency_key),
  164. CHECK (request_hash ~ '^[0-9a-f]{64}$'),
  165. CHECK (
  166. (status = 'completed' AND result IS NOT NULL
  167. AND result_hash ~ '^[0-9a-f]{64}$')
  168. OR status <> 'completed'
  169. )
  170. );
  171. CREATE INDEX idx_dataflow_deployment_operation_recovery
  172. ON public.dataflow_deployment_operations(
  173. status, updated_at, deployment_id
  174. )
  175. WHERE status IN ('claimed','failed','unknown');
  176. CREATE TABLE public.dataflow_deployment_transitions (
  177. id UUID PRIMARY KEY,
  178. deployment_id UUID NOT NULL
  179. REFERENCES public.dataflow_deployments(id) ON DELETE RESTRICT,
  180. operation_id UUID
  181. REFERENCES public.dataflow_deployment_operations(id)
  182. ON DELETE RESTRICT,
  183. from_status VARCHAR(30),
  184. to_status VARCHAR(30) NOT NULL,
  185. actor_uid UUID NOT NULL REFERENCES public.users(id) ON DELETE RESTRICT,
  186. reason VARCHAR(1000) NOT NULL,
  187. correlation_id UUID NOT NULL,
  188. detail JSONB NOT NULL DEFAULT '{}'::jsonb,
  189. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
  190. );
  191. CREATE INDEX idx_dataflow_deployment_transition_history
  192. ON public.dataflow_deployment_transitions(
  193. deployment_id, created_at DESC
  194. );
  195. """
  196. )
  197. def downgrade() -> None:
  198. raise RuntimeError(
  199. "Data Factory deployment evidence is forward-only and cannot downgrade"
  200. )