20260719_90_workflow_migration_control.py 4.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122
  1. """Add V55 dual-run, reconciliation, and cutover control records."""
  2. from alembic import op
  3. revision = "20260719_90"
  4. down_revision = "20260719_80"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.execute(
  9. """
  10. CREATE TABLE public.workflow_migration_states (
  11. id UUID PRIMARY KEY,
  12. dataflow_uid UUID NOT NULL,
  13. environment VARCHAR(20) NOT NULL
  14. CHECK (environment IN ('development', 'test', 'production')),
  15. mode VARCHAR(50) NOT NULL CHECK (mode IN (
  16. 'n8n_primary',
  17. 'n8n_primary_kestra_shadow',
  18. 'kestra_primary_n8n_standby',
  19. 'kestra_primary'
  20. )),
  21. status VARCHAR(30) NOT NULL DEFAULT 'stable'
  22. CHECK (status IN (
  23. 'stable', 'transition_pending', 'blocked', 'retired'
  24. )),
  25. n8n_definition_id VARCHAR(255),
  26. kestra_definition_id VARCHAR(255),
  27. migration_batch VARCHAR(50),
  28. observation_until TIMESTAMPTZ,
  29. metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
  30. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  31. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  32. UNIQUE (dataflow_uid, environment)
  33. );
  34. CREATE INDEX idx_workflow_migration_mode
  35. ON public.workflow_migration_states(
  36. environment, mode, status, updated_at DESC
  37. );
  38. CREATE TABLE public.workflow_dual_runs (
  39. id UUID PRIMARY KEY,
  40. migration_state_id UUID NOT NULL
  41. REFERENCES public.workflow_migration_states(id)
  42. ON DELETE CASCADE,
  43. input_hash CHAR(64) NOT NULL,
  44. n8n_execution_id VARCHAR(255),
  45. kestra_execution_id VARCHAR(255),
  46. n8n_status VARCHAR(30) NOT NULL,
  47. kestra_status VARCHAR(30) NOT NULL,
  48. isolation_mode VARCHAR(30) NOT NULL
  49. CHECK (isolation_mode IN ('read_only', 'isolated_target')),
  50. formal_engine_changed BOOLEAN NOT NULL DEFAULT FALSE,
  51. correlation_id UUID NOT NULL,
  52. safe_metrics JSONB NOT NULL DEFAULT '{}'::jsonb,
  53. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
  54. );
  55. CREATE INDEX idx_workflow_dual_runs_state_created
  56. ON public.workflow_dual_runs(
  57. migration_state_id, created_at DESC
  58. );
  59. CREATE TABLE public.workflow_reconciliation_reports (
  60. id UUID PRIMARY KEY,
  61. migration_state_id UUID NOT NULL
  62. REFERENCES public.workflow_migration_states(id)
  63. ON DELETE CASCADE,
  64. dual_run_id UUID
  65. REFERENCES public.workflow_dual_runs(id)
  66. ON DELETE SET NULL,
  67. status VARCHAR(20) NOT NULL
  68. CHECK (status IN ('passed', 'failed', 'blocked')),
  69. policy JSONB NOT NULL DEFAULT '{}'::jsonb,
  70. safe_metrics JSONB NOT NULL DEFAULT '{}'::jsonb,
  71. differences JSONB NOT NULL DEFAULT '[]'::jsonb,
  72. correlation_id UUID NOT NULL,
  73. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
  74. );
  75. CREATE INDEX idx_workflow_reconciliation_state_created
  76. ON public.workflow_reconciliation_reports(
  77. migration_state_id, status, created_at DESC
  78. );
  79. CREATE TABLE public.workflow_cutover_operations (
  80. id UUID PRIMARY KEY,
  81. migration_state_id UUID NOT NULL
  82. REFERENCES public.workflow_migration_states(id)
  83. ON DELETE CASCADE,
  84. idempotency_key VARCHAR(500) NOT NULL,
  85. source_mode VARCHAR(50) NOT NULL,
  86. target_mode VARCHAR(50) NOT NULL,
  87. status VARCHAR(30) NOT NULL
  88. CHECK (status IN (
  89. 'claimed', 'succeeded', 'rolled_back', 'failed'
  90. )),
  91. gates JSONB NOT NULL DEFAULT '{}'::jsonb,
  92. result JSONB NOT NULL DEFAULT '{}'::jsonb,
  93. safe_error VARCHAR(1000),
  94. correlation_id UUID NOT NULL,
  95. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  96. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  97. UNIQUE (migration_state_id, idempotency_key)
  98. );
  99. CREATE INDEX idx_workflow_cutover_state_created
  100. ON public.workflow_cutover_operations(
  101. migration_state_id, status, created_at DESC
  102. );
  103. """
  104. )
  105. def downgrade() -> None:
  106. op.execute(
  107. """
  108. DROP TABLE IF EXISTS public.workflow_cutover_operations;
  109. DROP TABLE IF EXISTS public.workflow_reconciliation_reports;
  110. DROP TABLE IF EXISTS public.workflow_dual_runs;
  111. DROP TABLE IF EXISTS public.workflow_migration_states;
  112. """
  113. )