"""Make governed DataFlow creation recoverable and idempotent.""" from alembic import op revision = "20260724_220" down_revision = "20260723_210" branch_labels = None depends_on = None def upgrade() -> None: op.execute( """ ALTER TABLE public.dataflow_draft_reservations ADD COLUMN state VARCHAR(20) NOT NULL DEFAULT 'reserved', ADD COLUMN lease_token UUID, ADD COLUMN lease_expires_at TIMESTAMPTZ, ADD COLUMN attempt_count INTEGER NOT NULL DEFAULT 0, ADD COLUMN dataflow_node_id BIGINT, ADD COLUMN result JSONB, ADD COLUMN result_digest CHAR(64), ADD COLUMN error_code VARCHAR(100), ADD COLUMN create_started_at TIMESTAMPTZ, ADD COLUMN completed_at TIMESTAMPTZ, ADD COLUMN failed_at TIMESTAMPTZ, ADD COLUMN updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP; UPDATE public.dataflow_draft_reservations SET state = 'failed', error_code = 'legacy_consumed_without_result', failed_at = COALESCE(consumed_at, CURRENT_TIMESTAMP) WHERE consumed_at IS NOT NULL; ALTER TABLE public.dataflow_draft_reservations ADD CONSTRAINT ck_dataflow_draft_saga_state CHECK (state IN ('reserved', 'creating', 'completed', 'failed')), ADD CONSTRAINT ck_dataflow_draft_saga_attempts CHECK (attempt_count >= 0), ADD CONSTRAINT ck_dataflow_draft_completed_result CHECK ( state <> 'completed' OR ( result IS NOT NULL AND result_digest IS NOT NULL AND dataflow_node_id IS NOT NULL AND completed_at IS NOT NULL ) ); CREATE INDEX idx_dataflow_draft_saga_claim ON public.dataflow_draft_reservations( state, lease_expires_at, expires_at ); """ ) def downgrade() -> None: raise RuntimeError( "DataFlow create saga reservations are forward-only and cannot downgrade" )