"""Add engine-neutral workflow definitions, bindings, schedules, and run audits.""" from alembic import op revision = "20260718_60" down_revision = "20260718_50" branch_labels = None depends_on = None def upgrade() -> None: op.execute( """ ALTER TABLE public.dataflow_workflow_versions ADD COLUMN engine_type VARCHAR(20) NOT NULL DEFAULT 'n8n', ADD COLUMN engine_definition_id VARCHAR(255), ADD COLUMN engine_revision VARCHAR(255), ADD COLUMN deployment_metadata JSONB NOT NULL DEFAULT '{}'::jsonb, ADD COLUMN workflow_spec JSONB, ADD COLUMN schedule_plan JSONB; UPDATE public.dataflow_workflow_versions SET engine_type = 'n8n', engine_definition_id = n8n_workflow_id, deployment_metadata = COALESCE(deployment_metadata, '{}'::jsonb) WHERE engine_definition_id IS NULL; ALTER TABLE public.dataflow_workflow_versions ALTER COLUMN engine_definition_id SET NOT NULL, ALTER COLUMN n8n_workflow_id DROP NOT NULL, ADD CONSTRAINT ck_workflow_version_engine_type CHECK (engine_type IN ('n8n', 'kestra')), ADD CONSTRAINT uq_workflow_version_identity UNIQUE (id, dataflow_uid, environment); CREATE INDEX idx_workflow_versions_engine_definition ON public.dataflow_workflow_versions( engine_type, engine_definition_id ); CREATE TABLE public.workflow_schedules ( id UUID PRIMARY KEY, workflow_version_id UUID NOT NULL REFERENCES public.dataflow_workflow_versions(id) ON DELETE CASCADE, schedule_plan JSONB NOT NULL, schedule_hash CHAR(64) NOT NULL, timezone VARCHAR(100) NOT NULL, status VARCHAR(20) NOT NULL DEFAULT 'draft' CHECK (status IN ('draft', 'active', 'paused', 'retired')), created_by UUID REFERENCES public.users(id) ON DELETE SET NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE (workflow_version_id, schedule_hash) ); CREATE INDEX idx_workflow_schedules_version ON public.workflow_schedules(workflow_version_id, status); CREATE TABLE public.workflow_engine_bindings ( id UUID PRIMARY KEY, workflow_version_id UUID NOT NULL, dataflow_uid UUID NOT NULL, environment VARCHAR(20) NOT NULL CHECK (environment IN ('development', 'test', 'production')), engine_type VARCHAR(20) NOT NULL CHECK (engine_type IN ('n8n', 'kestra')), engine_definition_id VARCHAR(255) NOT NULL, engine_revision VARCHAR(255), role VARCHAR(20) NOT NULL CHECK (role IN ('primary', 'shadow', 'standby', 'archived')), status VARCHAR(20) NOT NULL DEFAULT 'disabled' CHECK (status IN ('enabled', 'disabled')), metadata JSONB NOT NULL DEFAULT '{}'::jsonb, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, FOREIGN KEY (workflow_version_id, dataflow_uid, environment) REFERENCES public.dataflow_workflow_versions( id, dataflow_uid, environment ) ON DELETE CASCADE, UNIQUE (engine_type, engine_definition_id) ); CREATE UNIQUE INDEX uq_workflow_engine_primary_environment ON public.workflow_engine_bindings(dataflow_uid, environment) WHERE role = 'primary' AND status = 'enabled'; CREATE UNIQUE INDEX uq_workflow_engine_shadow_environment ON public.workflow_engine_bindings(dataflow_uid, environment) WHERE role = 'shadow' AND status = 'enabled'; CREATE INDEX idx_workflow_engine_binding_version ON public.workflow_engine_bindings(workflow_version_id, role, status); CREATE TABLE public.workflow_runs ( id UUID PRIMARY KEY, workflow_version_id UUID NOT NULL REFERENCES public.dataflow_workflow_versions(id) ON DELETE RESTRICT, schedule_id UUID REFERENCES public.workflow_schedules(id) ON DELETE SET NULL, engine_type VARCHAR(20) NOT NULL CHECK (engine_type IN ('n8n', 'kestra')), engine_execution_id VARCHAR(255), trigger_type VARCHAR(20) NOT NULL CHECK (trigger_type IN ( 'manual', 'cron', 'at', 'event', 'backfill', 'replay' )), status VARCHAR(30) NOT NULL CHECK (status IN ( 'queued', 'running', 'success', 'failed', 'paused', 'killed', 'cancelled', 'unknown' )), correlation_id UUID NOT NULL, idempotency_key VARCHAR(500), started_at TIMESTAMPTZ, finished_at TIMESTAMPTZ, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP ); CREATE UNIQUE INDEX uq_workflow_run_engine_execution ON public.workflow_runs(engine_type, engine_execution_id) WHERE engine_execution_id IS NOT NULL; CREATE INDEX idx_workflow_runs_version_created ON public.workflow_runs(workflow_version_id, created_at DESC); CREATE INDEX idx_workflow_runs_correlation ON public.workflow_runs(correlation_id); CREATE TABLE public.workflow_task_runs ( id UUID PRIMARY KEY, workflow_run_id UUID NOT NULL REFERENCES public.workflow_runs(id) ON DELETE CASCADE, node_id VARCHAR(100) NOT NULL, engine_task_id VARCHAR(255), status VARCHAR(30) NOT NULL CHECK (status IN ( 'queued', 'running', 'success', 'failed', 'paused', 'killed', 'cancelled', 'unknown' )), attempt INTEGER NOT NULL DEFAULT 1 CHECK (attempt > 0), idempotency_key VARCHAR(500), safe_error VARCHAR(1000), started_at TIMESTAMPTZ, finished_at TIMESTAMPTZ, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE (workflow_run_id, node_id, attempt) ); CREATE INDEX idx_workflow_task_runs_run ON public.workflow_task_runs(workflow_run_id, node_id); CREATE TABLE public.workflow_plan_audits ( id UUID PRIMARY KEY, workflow_version_id UUID REFERENCES public.dataflow_workflow_versions(id) ON DELETE SET NULL, actor_uid UUID REFERENCES public.users(id) ON DELETE SET NULL, action VARCHAR(80) NOT NULL, model_provider VARCHAR(80), model_name VARCHAR(120), prompt_version VARCHAR(80), schema_version VARCHAR(40) NOT NULL, context_hash CHAR(64), candidate_hash CHAR(64), decision VARCHAR(30) NOT NULL CHECK (decision IN ('allowed', 'rejected', 'failed', 'recorded')), decision_detail JSONB NOT NULL DEFAULT '{}'::jsonb, correlation_id UUID NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX idx_workflow_plan_audits_correlation ON public.workflow_plan_audits(correlation_id, created_at); CREATE INDEX idx_workflow_plan_audits_version ON public.workflow_plan_audits( workflow_version_id, created_at DESC ); """ ) def downgrade() -> None: op.execute( """ DROP TABLE IF EXISTS public.workflow_plan_audits; DROP TABLE IF EXISTS public.workflow_task_runs; DROP TABLE IF EXISTS public.workflow_runs; DROP TABLE IF EXISTS public.workflow_engine_bindings; DROP TABLE IF EXISTS public.workflow_schedules; DROP INDEX IF EXISTS public.idx_workflow_versions_engine_definition; DO $$ BEGIN IF EXISTS ( SELECT 1 FROM public.dataflow_workflow_versions WHERE n8n_workflow_id IS NULL ) THEN RAISE EXCEPTION 'cannot downgrade while non-n8n workflow versions exist'; END IF; END $$; ALTER TABLE public.dataflow_workflow_versions ALTER COLUMN n8n_workflow_id SET NOT NULL, DROP CONSTRAINT IF EXISTS uq_workflow_version_identity, DROP CONSTRAINT IF EXISTS ck_workflow_version_engine_type, DROP COLUMN IF EXISTS schedule_plan, DROP COLUMN IF EXISTS workflow_spec, DROP COLUMN IF EXISTS deployment_metadata, DROP COLUMN IF EXISTS engine_revision, DROP COLUMN IF EXISTS engine_definition_id, DROP COLUMN IF EXISTS engine_type; """ )