20260723_170_rule_execution_attempts.py 2.8 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374
  1. """Add exact task identity, leases, replay, and SQL staging receipts."""
  2. from alembic import op
  3. revision = "20260723_170"
  4. down_revision = "20260723_160"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.execute(
  9. """
  10. ALTER TABLE public.rule_runs
  11. ADD COLUMN attempt_no INTEGER NOT NULL DEFAULT 1
  12. CHECK (attempt_no > 0),
  13. ADD COLUMN lease_owner UUID,
  14. ADD COLUMN lease_expires_at TIMESTAMPTZ,
  15. ADD COLUMN heartbeat_at TIMESTAMPTZ,
  16. ADD COLUMN evidence_digest CHAR(64);
  17. CREATE INDEX idx_rule_runs_lease
  18. ON public.rule_runs(status, lease_expires_at, id);
  19. ALTER TABLE public.runner_task_executions
  20. ADD COLUMN deployment_id UUID,
  21. ADD COLUMN environment VARCHAR(20)
  22. CHECK (environment IN (
  23. 'development','test','production'
  24. )),
  25. ADD COLUMN replay_http_status SMALLINT
  26. CHECK (replay_http_status BETWEEN 200 AND 599),
  27. ADD COLUMN replay_body JSONB,
  28. ADD COLUMN replay_digest CHAR(64);
  29. CREATE INDEX idx_runner_task_execution_identity
  30. ON public.runner_task_executions
  31. (deployment_id, environment, correlation_id, node_id);
  32. CREATE TABLE public.rule_sql_staging_receipts (
  33. id UUID PRIMARY KEY,
  34. producer_rule_run_id UUID NOT NULL
  35. REFERENCES public.rule_runs(id) ON DELETE RESTRICT,
  36. deployment_id UUID NOT NULL
  37. REFERENCES public.dataflow_deployments(id)
  38. ON DELETE RESTRICT,
  39. correlation_id UUID NOT NULL,
  40. output_binding_id UUID NOT NULL
  41. REFERENCES public.dataflow_dataset_bindings(id)
  42. ON DELETE RESTRICT,
  43. output_binding_hash CHAR(64) NOT NULL,
  44. relation_ref VARCHAR(1000) NOT NULL,
  45. relation_digest CHAR(64) NOT NULL,
  46. commit_outcome VARCHAR(30) NOT NULL
  47. CHECK (commit_outcome IN ('committed','unknown')),
  48. status VARCHAR(20) NOT NULL
  49. CHECK (status IN ('pending','ready','failed','expired')),
  50. expires_at TIMESTAMPTZ NOT NULL,
  51. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  52. ready_at TIMESTAMPTZ,
  53. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  54. UNIQUE (producer_rule_run_id, output_binding_id)
  55. );
  56. CREATE INDEX idx_rule_sql_staging_receipt_lookup
  57. ON public.rule_sql_staging_receipts
  58. (id, deployment_id, correlation_id, status, expires_at);
  59. """
  60. )
  61. def downgrade() -> None:
  62. raise RuntimeError(
  63. "rule attempt leases and staging receipts are forward-only"
  64. )