| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211 |
- """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;
- """
- )
|