"""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" )