| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374 |
- """Add exact task identity, leases, replay, and SQL staging receipts."""
- from alembic import op
- revision = "20260723_170"
- down_revision = "20260723_160"
- branch_labels = None
- depends_on = None
- def upgrade() -> None:
- op.execute(
- """
- ALTER TABLE public.rule_runs
- ADD COLUMN attempt_no INTEGER NOT NULL DEFAULT 1
- CHECK (attempt_no > 0),
- ADD COLUMN lease_owner UUID,
- ADD COLUMN lease_expires_at TIMESTAMPTZ,
- ADD COLUMN heartbeat_at TIMESTAMPTZ,
- ADD COLUMN evidence_digest CHAR(64);
- CREATE INDEX idx_rule_runs_lease
- ON public.rule_runs(status, lease_expires_at, id);
- ALTER TABLE public.runner_task_executions
- ADD COLUMN deployment_id UUID,
- ADD COLUMN environment VARCHAR(20)
- CHECK (environment IN (
- 'development','test','production'
- )),
- ADD COLUMN replay_http_status SMALLINT
- CHECK (replay_http_status BETWEEN 200 AND 599),
- ADD COLUMN replay_body JSONB,
- ADD COLUMN replay_digest CHAR(64);
- CREATE INDEX idx_runner_task_execution_identity
- ON public.runner_task_executions
- (deployment_id, environment, correlation_id, node_id);
- CREATE TABLE public.rule_sql_staging_receipts (
- id UUID PRIMARY KEY,
- producer_rule_run_id UUID NOT NULL
- REFERENCES public.rule_runs(id) ON DELETE RESTRICT,
- deployment_id UUID NOT NULL
- REFERENCES public.dataflow_deployments(id)
- ON DELETE RESTRICT,
- correlation_id UUID NOT NULL,
- output_binding_id UUID NOT NULL
- REFERENCES public.dataflow_dataset_bindings(id)
- ON DELETE RESTRICT,
- output_binding_hash CHAR(64) NOT NULL,
- relation_ref VARCHAR(1000) NOT NULL,
- relation_digest CHAR(64) NOT NULL,
- commit_outcome VARCHAR(30) NOT NULL
- CHECK (commit_outcome IN ('committed','unknown')),
- status VARCHAR(20) NOT NULL
- CHECK (status IN ('pending','ready','failed','expired')),
- expires_at TIMESTAMPTZ NOT NULL,
- created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
- ready_at TIMESTAMPTZ,
- updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
- UNIQUE (producer_rule_run_id, output_binding_id)
- );
- CREATE INDEX idx_rule_sql_staging_receipt_lookup
- ON public.rule_sql_staging_receipts
- (id, deployment_id, correlation_id, status, expires_at);
- """
- )
- def downgrade() -> None:
- raise RuntimeError(
- "rule attempt leases and staging receipts are forward-only"
- )
|