20260718_60_workflow_engine_abstraction.py 9.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211
  1. """Add engine-neutral workflow definitions, bindings, schedules, and run audits."""
  2. from alembic import op
  3. revision = "20260718_60"
  4. down_revision = "20260718_50"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.execute(
  9. """
  10. ALTER TABLE public.dataflow_workflow_versions
  11. ADD COLUMN engine_type VARCHAR(20) NOT NULL DEFAULT 'n8n',
  12. ADD COLUMN engine_definition_id VARCHAR(255),
  13. ADD COLUMN engine_revision VARCHAR(255),
  14. ADD COLUMN deployment_metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
  15. ADD COLUMN workflow_spec JSONB,
  16. ADD COLUMN schedule_plan JSONB;
  17. UPDATE public.dataflow_workflow_versions
  18. SET engine_type = 'n8n',
  19. engine_definition_id = n8n_workflow_id,
  20. deployment_metadata = COALESCE(deployment_metadata, '{}'::jsonb)
  21. WHERE engine_definition_id IS NULL;
  22. ALTER TABLE public.dataflow_workflow_versions
  23. ALTER COLUMN engine_definition_id SET NOT NULL,
  24. ALTER COLUMN n8n_workflow_id DROP NOT NULL,
  25. ADD CONSTRAINT ck_workflow_version_engine_type
  26. CHECK (engine_type IN ('n8n', 'kestra')),
  27. ADD CONSTRAINT uq_workflow_version_identity
  28. UNIQUE (id, dataflow_uid, environment);
  29. CREATE INDEX idx_workflow_versions_engine_definition
  30. ON public.dataflow_workflow_versions(
  31. engine_type, engine_definition_id
  32. );
  33. CREATE TABLE public.workflow_schedules (
  34. id UUID PRIMARY KEY,
  35. workflow_version_id UUID NOT NULL
  36. REFERENCES public.dataflow_workflow_versions(id)
  37. ON DELETE CASCADE,
  38. schedule_plan JSONB NOT NULL,
  39. schedule_hash CHAR(64) NOT NULL,
  40. timezone VARCHAR(100) NOT NULL,
  41. status VARCHAR(20) NOT NULL DEFAULT 'draft'
  42. CHECK (status IN ('draft', 'active', 'paused', 'retired')),
  43. created_by UUID REFERENCES public.users(id) ON DELETE SET NULL,
  44. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  45. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  46. UNIQUE (workflow_version_id, schedule_hash)
  47. );
  48. CREATE INDEX idx_workflow_schedules_version
  49. ON public.workflow_schedules(workflow_version_id, status);
  50. CREATE TABLE public.workflow_engine_bindings (
  51. id UUID PRIMARY KEY,
  52. workflow_version_id UUID NOT NULL,
  53. dataflow_uid UUID NOT NULL,
  54. environment VARCHAR(20) NOT NULL
  55. CHECK (environment IN ('development', 'test', 'production')),
  56. engine_type VARCHAR(20) NOT NULL
  57. CHECK (engine_type IN ('n8n', 'kestra')),
  58. engine_definition_id VARCHAR(255) NOT NULL,
  59. engine_revision VARCHAR(255),
  60. role VARCHAR(20) NOT NULL
  61. CHECK (role IN ('primary', 'shadow', 'standby', 'archived')),
  62. status VARCHAR(20) NOT NULL DEFAULT 'disabled'
  63. CHECK (status IN ('enabled', 'disabled')),
  64. metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
  65. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  66. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  67. FOREIGN KEY (workflow_version_id, dataflow_uid, environment)
  68. REFERENCES public.dataflow_workflow_versions(
  69. id, dataflow_uid, environment
  70. ) ON DELETE CASCADE,
  71. UNIQUE (engine_type, engine_definition_id)
  72. );
  73. CREATE UNIQUE INDEX uq_workflow_engine_primary_environment
  74. ON public.workflow_engine_bindings(dataflow_uid, environment)
  75. WHERE role = 'primary' AND status = 'enabled';
  76. CREATE UNIQUE INDEX uq_workflow_engine_shadow_environment
  77. ON public.workflow_engine_bindings(dataflow_uid, environment)
  78. WHERE role = 'shadow' AND status = 'enabled';
  79. CREATE INDEX idx_workflow_engine_binding_version
  80. ON public.workflow_engine_bindings(workflow_version_id, role, status);
  81. CREATE TABLE public.workflow_runs (
  82. id UUID PRIMARY KEY,
  83. workflow_version_id UUID NOT NULL
  84. REFERENCES public.dataflow_workflow_versions(id)
  85. ON DELETE RESTRICT,
  86. schedule_id UUID REFERENCES public.workflow_schedules(id)
  87. ON DELETE SET NULL,
  88. engine_type VARCHAR(20) NOT NULL
  89. CHECK (engine_type IN ('n8n', 'kestra')),
  90. engine_execution_id VARCHAR(255),
  91. trigger_type VARCHAR(20) NOT NULL
  92. CHECK (trigger_type IN (
  93. 'manual', 'cron', 'at', 'event', 'backfill', 'replay'
  94. )),
  95. status VARCHAR(30) NOT NULL
  96. CHECK (status IN (
  97. 'queued', 'running', 'success', 'failed', 'paused',
  98. 'killed', 'cancelled', 'unknown'
  99. )),
  100. correlation_id UUID NOT NULL,
  101. idempotency_key VARCHAR(500),
  102. started_at TIMESTAMPTZ,
  103. finished_at TIMESTAMPTZ,
  104. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
  105. );
  106. CREATE UNIQUE INDEX uq_workflow_run_engine_execution
  107. ON public.workflow_runs(engine_type, engine_execution_id)
  108. WHERE engine_execution_id IS NOT NULL;
  109. CREATE INDEX idx_workflow_runs_version_created
  110. ON public.workflow_runs(workflow_version_id, created_at DESC);
  111. CREATE INDEX idx_workflow_runs_correlation
  112. ON public.workflow_runs(correlation_id);
  113. CREATE TABLE public.workflow_task_runs (
  114. id UUID PRIMARY KEY,
  115. workflow_run_id UUID NOT NULL
  116. REFERENCES public.workflow_runs(id) ON DELETE CASCADE,
  117. node_id VARCHAR(100) NOT NULL,
  118. engine_task_id VARCHAR(255),
  119. status VARCHAR(30) NOT NULL
  120. CHECK (status IN (
  121. 'queued', 'running', 'success', 'failed', 'paused',
  122. 'killed', 'cancelled', 'unknown'
  123. )),
  124. attempt INTEGER NOT NULL DEFAULT 1 CHECK (attempt > 0),
  125. idempotency_key VARCHAR(500),
  126. safe_error VARCHAR(1000),
  127. started_at TIMESTAMPTZ,
  128. finished_at TIMESTAMPTZ,
  129. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  130. UNIQUE (workflow_run_id, node_id, attempt)
  131. );
  132. CREATE INDEX idx_workflow_task_runs_run
  133. ON public.workflow_task_runs(workflow_run_id, node_id);
  134. CREATE TABLE public.workflow_plan_audits (
  135. id UUID PRIMARY KEY,
  136. workflow_version_id UUID
  137. REFERENCES public.dataflow_workflow_versions(id)
  138. ON DELETE SET NULL,
  139. actor_uid UUID REFERENCES public.users(id) ON DELETE SET NULL,
  140. action VARCHAR(80) NOT NULL,
  141. model_provider VARCHAR(80),
  142. model_name VARCHAR(120),
  143. prompt_version VARCHAR(80),
  144. schema_version VARCHAR(40) NOT NULL,
  145. context_hash CHAR(64),
  146. candidate_hash CHAR(64),
  147. decision VARCHAR(30) NOT NULL
  148. CHECK (decision IN ('allowed', 'rejected', 'failed', 'recorded')),
  149. decision_detail JSONB NOT NULL DEFAULT '{}'::jsonb,
  150. correlation_id UUID NOT NULL,
  151. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
  152. );
  153. CREATE INDEX idx_workflow_plan_audits_correlation
  154. ON public.workflow_plan_audits(correlation_id, created_at);
  155. CREATE INDEX idx_workflow_plan_audits_version
  156. ON public.workflow_plan_audits(
  157. workflow_version_id, created_at DESC
  158. );
  159. """
  160. )
  161. def downgrade() -> None:
  162. op.execute(
  163. """
  164. DROP TABLE IF EXISTS public.workflow_plan_audits;
  165. DROP TABLE IF EXISTS public.workflow_task_runs;
  166. DROP TABLE IF EXISTS public.workflow_runs;
  167. DROP TABLE IF EXISTS public.workflow_engine_bindings;
  168. DROP TABLE IF EXISTS public.workflow_schedules;
  169. DROP INDEX IF EXISTS public.idx_workflow_versions_engine_definition;
  170. DO $$
  171. BEGIN
  172. IF EXISTS (
  173. SELECT 1
  174. FROM public.dataflow_workflow_versions
  175. WHERE n8n_workflow_id IS NULL
  176. ) THEN
  177. RAISE EXCEPTION
  178. 'cannot downgrade while non-n8n workflow versions exist';
  179. END IF;
  180. END
  181. $$;
  182. ALTER TABLE public.dataflow_workflow_versions
  183. ALTER COLUMN n8n_workflow_id SET NOT NULL,
  184. DROP CONSTRAINT IF EXISTS uq_workflow_version_identity,
  185. DROP CONSTRAINT IF EXISTS ck_workflow_version_engine_type,
  186. DROP COLUMN IF EXISTS schedule_plan,
  187. DROP COLUMN IF EXISTS workflow_spec,
  188. DROP COLUMN IF EXISTS deployment_metadata,
  189. DROP COLUMN IF EXISTS engine_revision,
  190. DROP COLUMN IF EXISTS engine_definition_id,
  191. DROP COLUMN IF EXISTS engine_type;
  192. """
  193. )