| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214 |
- """Add recoverable immutable Data Factory deployment lifecycle."""
- from alembic import op
- revision = "20260724_240"
- down_revision = "20260724_230"
- branch_labels = None
- depends_on = None
- def upgrade() -> None:
- op.execute(
- """
- ALTER TABLE public.dataflow_deployments
- ADD COLUMN dataflow_uid UUID,
- ADD COLUMN package JSONB,
- ADD COLUMN package_hash CHAR(64),
- ADD COLUMN standard_version_ids JSONB NOT NULL DEFAULT '[]'::jsonb,
- ADD COLUMN rule_version_ids JSONB NOT NULL DEFAULT '[]'::jsonb,
- ADD COLUMN binding_snapshot JSONB,
- ADD COLUMN binding_hash CHAR(64),
- ADD COLUMN schema_snapshots JSONB,
- ADD COLUMN schema_snapshot_hash CHAR(64),
- ADD COLUMN physical_plan_hashes JSONB NOT NULL DEFAULT '[]'::jsonb,
- ADD COLUMN workflow_spec JSONB,
- ADD COLUMN workflow_spec_hash CHAR(64),
- ADD COLUMN schedule_snapshot JSONB,
- ADD COLUMN schedule_hash CHAR(64),
- ADD COLUMN engine_namespace VARCHAR(200),
- ADD COLUMN engine_definition_id VARCHAR(200),
- ADD COLUMN engine_revision VARCHAR(200),
- ADD COLUMN engine_definition_hash CHAR(64),
- ADD COLUMN deploy_receipt JSONB NOT NULL DEFAULT '{}'::jsonb,
- ADD COLUMN canary_evidence_id UUID,
- ADD COLUMN previous_active_deployment_id UUID
- REFERENCES public.dataflow_deployments(id) ON DELETE RESTRICT,
- ADD COLUMN created_by UUID REFERENCES public.users(id) ON DELETE SET NULL,
- ADD COLUMN create_reason VARCHAR(1000),
- ADD COLUMN create_idempotency_key VARCHAR(200),
- ADD COLUMN correlation_id UUID,
- ADD COLUMN lock_version INTEGER NOT NULL DEFAULT 0,
- ADD COLUMN failure_code VARCHAR(100),
- ADD COLUMN superseded_at TIMESTAMPTZ,
- ADD COLUMN rolled_back_at TIMESTAMPTZ;
- UPDATE public.dataflow_deployments d
- SET dataflow_uid = v.dataflow_uid
- FROM public.dataflow_versions v
- WHERE v.id = d.dataflow_version_id
- AND d.dataflow_uid IS NULL;
- ALTER TABLE public.dataflow_deployments
- ADD CONSTRAINT ck_dataflow_deployment_lock_version
- CHECK (lock_version >= 0),
- ADD CONSTRAINT ck_dataflow_deployment_hash_shape
- CHECK (
- (package_hash IS NULL OR package_hash ~ '^[0-9a-f]{64}$')
- AND (binding_hash IS NULL OR binding_hash ~ '^[0-9a-f]{64}$')
- AND (
- schema_snapshot_hash IS NULL
- OR schema_snapshot_hash ~ '^[0-9a-f]{64}$'
- )
- AND (
- workflow_spec_hash IS NULL
- OR workflow_spec_hash ~ '^[0-9a-f]{64}$'
- )
- AND (
- schedule_hash IS NULL
- OR schedule_hash ~ '^[0-9a-f]{64}$'
- )
- AND (
- engine_definition_hash IS NULL
- OR engine_definition_hash ~ '^[0-9a-f]{64}$'
- )
- ),
- ADD CONSTRAINT ck_dataflow_deployment_snapshot_complete
- CHECK (
- package IS NULL
- OR (
- dataflow_uid IS NOT NULL
- AND package_hash IS NOT NULL
- AND binding_snapshot IS NOT NULL
- AND binding_hash IS NOT NULL
- AND schema_snapshots IS NOT NULL
- AND schema_snapshot_hash IS NOT NULL
- AND jsonb_array_length(physical_plan_hashes) > 0
- AND workflow_spec IS NOT NULL
- AND workflow_spec_hash IS NOT NULL
- AND schedule_snapshot IS NOT NULL
- AND schedule_hash IS NOT NULL
- AND created_by IS NOT NULL
- AND create_reason IS NOT NULL
- AND create_idempotency_key IS NOT NULL
- AND correlation_id IS NOT NULL
- )
- );
- ALTER TABLE public.dataflow_deployments
- DROP CONSTRAINT IF EXISTS
- dataflow_deployments_dataflow_version_id_environment_key;
- DROP INDEX IF EXISTS public.uq_dataflow_deployment_active_environment;
- CREATE UNIQUE INDEX uq_dataflow_deployment_active_environment
- ON public.dataflow_deployments(dataflow_uid, environment)
- WHERE status = 'active';
- CREATE UNIQUE INDEX uq_dataflow_deployment_create_idempotency
- ON public.dataflow_deployments(
- created_by, environment, create_idempotency_key
- )
- WHERE create_idempotency_key IS NOT NULL;
- CREATE INDEX idx_dataflow_deployment_factory_list
- ON public.dataflow_deployments(
- environment, updated_at DESC, dataflow_uid
- )
- WHERE package IS NOT NULL;
- CREATE TABLE public.dataflow_canary_evidence (
- id UUID PRIMARY KEY,
- deployment_id UUID NOT NULL
- REFERENCES public.dataflow_deployments(id) ON DELETE RESTRICT,
- execution_id VARCHAR(200) NOT NULL,
- status VARCHAR(20) NOT NULL CHECK (
- status IN ('started','passed','failed','unknown')
- ),
- package_hash CHAR(64) NOT NULL,
- binding_hash CHAR(64) NOT NULL,
- schema_snapshot_hash CHAR(64) NOT NULL,
- physical_plan_hashes JSONB NOT NULL,
- workflow_spec_hash CHAR(64) NOT NULL,
- schedule_hash CHAR(64) NOT NULL,
- engine_definition_hash CHAR(64) NOT NULL,
- verified_by UUID NOT NULL REFERENCES public.users(id) ON DELETE RESTRICT,
- engine_response JSONB NOT NULL DEFAULT '{}'::jsonb,
- expires_at TIMESTAMPTZ NOT NULL,
- created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
- UNIQUE (deployment_id, execution_id)
- );
- CREATE INDEX idx_dataflow_canary_activation_gate
- ON public.dataflow_canary_evidence(
- deployment_id, status, expires_at DESC
- );
- ALTER TABLE public.dataflow_deployments
- ADD CONSTRAINT fk_dataflow_deployment_canary_evidence
- FOREIGN KEY (canary_evidence_id)
- REFERENCES public.dataflow_canary_evidence(id)
- ON DELETE RESTRICT;
- CREATE TABLE public.dataflow_deployment_operations (
- id UUID PRIMARY KEY,
- deployment_id UUID NOT NULL
- REFERENCES public.dataflow_deployments(id) ON DELETE RESTRICT,
- action VARCHAR(30) NOT NULL CHECK (
- action IN (
- 'deploy_disabled','run_canary','activate','rollback'
- )
- ),
- idempotency_key VARCHAR(200) NOT NULL,
- status VARCHAR(20) NOT NULL CHECK (
- status IN ('claimed','completed','failed','unknown')
- ),
- actor_uid UUID NOT NULL REFERENCES public.users(id) ON DELETE RESTRICT,
- correlation_id UUID NOT NULL,
- reason VARCHAR(1000) NOT NULL,
- request JSONB NOT NULL DEFAULT '{}'::jsonb,
- request_hash CHAR(64) NOT NULL,
- result JSONB,
- result_hash CHAR(64),
- error_code VARCHAR(100),
- started_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
- completed_at TIMESTAMPTZ,
- failed_at TIMESTAMPTZ,
- updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
- UNIQUE (deployment_id, action, idempotency_key),
- CHECK (request_hash ~ '^[0-9a-f]{64}$'),
- CHECK (
- (status = 'completed' AND result IS NOT NULL
- AND result_hash ~ '^[0-9a-f]{64}$')
- OR status <> 'completed'
- )
- );
- CREATE INDEX idx_dataflow_deployment_operation_recovery
- ON public.dataflow_deployment_operations(
- status, updated_at, deployment_id
- )
- WHERE status IN ('claimed','failed','unknown');
- CREATE TABLE public.dataflow_deployment_transitions (
- id UUID PRIMARY KEY,
- deployment_id UUID NOT NULL
- REFERENCES public.dataflow_deployments(id) ON DELETE RESTRICT,
- operation_id UUID
- REFERENCES public.dataflow_deployment_operations(id)
- ON DELETE RESTRICT,
- from_status VARCHAR(30),
- to_status VARCHAR(30) NOT NULL,
- actor_uid UUID NOT NULL REFERENCES public.users(id) ON DELETE RESTRICT,
- reason VARCHAR(1000) NOT NULL,
- correlation_id UUID NOT NULL,
- detail JSONB NOT NULL DEFAULT '{}'::jsonb,
- created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
- );
- CREATE INDEX idx_dataflow_deployment_transition_history
- ON public.dataflow_deployment_transitions(
- deployment_id, created_at DESC
- );
- """
- )
- def downgrade() -> None:
- raise RuntimeError(
- "Data Factory deployment evidence is forward-only and cannot downgrade"
- )
|