20260719_80_mcp_gateway_control_plane.py 5.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124
  1. """Add persistent MCP Gateway scope, canary evidence, and idempotency state."""
  2. from alembic import op
  3. revision = "20260719_80"
  4. down_revision = "20260719_70"
  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 business_domain VARCHAR(200) NOT NULL DEFAULT 'legacy',
  12. ADD COLUMN created_by_subject VARCHAR(200) NOT NULL DEFAULT 'legacy',
  13. ADD COLUMN write_authorized BOOLEAN NOT NULL DEFAULT FALSE;
  14. CREATE INDEX idx_workflow_versions_gateway_scope
  15. ON public.dataflow_workflow_versions(
  16. business_domain, environment, status, version_no DESC
  17. );
  18. ALTER TABLE public.workflow_plan_audits
  19. ADD COLUMN actor_subject VARCHAR(200) NOT NULL DEFAULT 'legacy';
  20. ALTER TABLE public.workflow_plan_audits
  21. DROP CONSTRAINT IF EXISTS workflow_plan_audits_decision_check;
  22. ALTER TABLE public.workflow_plan_audits
  23. DROP CONSTRAINT IF EXISTS ck_workflow_plan_audits_decision;
  24. ALTER TABLE public.workflow_plan_audits
  25. ADD CONSTRAINT ck_workflow_plan_audits_decision
  26. CHECK (decision IN (
  27. 'allowed', 'rejected', 'failed', 'recorded',
  28. 'idempotent_replay'
  29. ));
  30. CREATE TABLE public.workflow_canary_evidence (
  31. id UUID PRIMARY KEY,
  32. workflow_version_id UUID NOT NULL
  33. REFERENCES public.dataflow_workflow_versions(id)
  34. ON DELETE CASCADE,
  35. engine_execution_id VARCHAR(255) NOT NULL,
  36. baseline_execution_id VARCHAR(255),
  37. status VARCHAR(20) NOT NULL
  38. CHECK (status IN ('started', 'passed', 'failed')),
  39. sample_runs INTEGER NOT NULL DEFAULT 0 CHECK (sample_runs >= 0),
  40. verified_by VARCHAR(200),
  41. verification_summary JSONB NOT NULL DEFAULT '{}'::jsonb,
  42. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  43. verified_at TIMESTAMPTZ,
  44. UNIQUE (engine_execution_id)
  45. );
  46. CREATE INDEX idx_workflow_canary_version_created
  47. ON public.workflow_canary_evidence(
  48. workflow_version_id, created_at DESC
  49. );
  50. CREATE TABLE public.workflow_gateway_operations (
  51. id UUID PRIMARY KEY,
  52. workflow_version_id UUID
  53. REFERENCES public.dataflow_workflow_versions(id)
  54. ON DELETE SET NULL,
  55. action VARCHAR(80) NOT NULL,
  56. idempotency_key VARCHAR(500) NOT NULL,
  57. status VARCHAR(20) NOT NULL
  58. CHECK (status IN ('claimed', 'succeeded', 'failed')),
  59. result JSONB NOT NULL DEFAULT '{}'::jsonb,
  60. safe_error VARCHAR(1000),
  61. actor_subject VARCHAR(200) NOT NULL,
  62. correlation_id UUID NOT NULL,
  63. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  64. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  65. UNIQUE (action, idempotency_key)
  66. );
  67. CREATE INDEX idx_workflow_gateway_operations_version
  68. ON public.workflow_gateway_operations(
  69. workflow_version_id, action, created_at DESC
  70. );
  71. ALTER TABLE public.workflow_runs
  72. ADD COLUMN parent_run_id UUID
  73. REFERENCES public.workflow_runs(id) ON DELETE SET NULL,
  74. ADD COLUMN attempt INTEGER NOT NULL DEFAULT 1 CHECK (attempt > 0),
  75. ADD COLUMN run_metadata JSONB NOT NULL DEFAULT '{}'::jsonb;
  76. CREATE INDEX idx_workflow_runs_parent
  77. ON public.workflow_runs(parent_run_id, attempt);
  78. """
  79. )
  80. def downgrade() -> None:
  81. op.execute(
  82. """
  83. DROP TABLE IF EXISTS public.workflow_gateway_operations;
  84. DROP TABLE IF EXISTS public.workflow_canary_evidence;
  85. DROP INDEX IF EXISTS public.idx_workflow_runs_parent;
  86. ALTER TABLE public.workflow_runs
  87. DROP COLUMN IF EXISTS run_metadata,
  88. DROP COLUMN IF EXISTS attempt,
  89. DROP COLUMN IF EXISTS parent_run_id;
  90. UPDATE public.workflow_plan_audits
  91. SET decision = 'recorded'
  92. WHERE decision = 'idempotent_replay';
  93. ALTER TABLE public.workflow_plan_audits
  94. DROP CONSTRAINT IF EXISTS ck_workflow_plan_audits_decision;
  95. ALTER TABLE public.workflow_plan_audits
  96. ADD CONSTRAINT workflow_plan_audits_decision_check
  97. CHECK (decision IN (
  98. 'allowed', 'rejected', 'failed', 'recorded'
  99. ));
  100. ALTER TABLE public.workflow_plan_audits
  101. DROP COLUMN IF EXISTS actor_subject;
  102. DROP INDEX IF EXISTS public.idx_workflow_versions_gateway_scope;
  103. ALTER TABLE public.dataflow_workflow_versions
  104. DROP COLUMN IF EXISTS write_authorized,
  105. DROP COLUMN IF EXISTS created_by_subject,
  106. DROP COLUMN IF EXISTS business_domain;
  107. """
  108. )