"""Create the transactional outbox and idempotent consumption ledger.""" from alembic import op revision = "20260716_03" down_revision = "20260716_02" branch_labels = None depends_on = None def upgrade() -> None: op.execute( """ CREATE TABLE IF NOT EXISTS public.outbox_events ( event_id UUID PRIMARY KEY, aggregate_type VARCHAR(100) NOT NULL, aggregate_id VARCHAR(200) NOT NULL, event_type VARCHAR(150) NOT NULL, payload JSONB NOT NULL, correlation_id VARCHAR(100), status VARCHAR(30) NOT NULL DEFAULT 'pending', attempts INTEGER NOT NULL DEFAULT 0, available_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, published_at TIMESTAMPTZ, last_error TEXT, CONSTRAINT ck_outbox_status CHECK ( status IN ('pending', 'processing', 'published', 'failed', 'dead_letter') ) ); CREATE INDEX IF NOT EXISTS idx_outbox_delivery ON public.outbox_events(status, available_at, created_at); CREATE INDEX IF NOT EXISTS idx_outbox_aggregate ON public.outbox_events(aggregate_type, aggregate_id); CREATE TABLE IF NOT EXISTS public.outbox_consumptions ( consumer_name VARCHAR(100) NOT NULL, event_id UUID NOT NULL REFERENCES public.outbox_events(event_id), processed_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (consumer_name, event_id) ); CREATE OR REPLACE FUNCTION public.touch_outbox_updated_at() RETURNS TRIGGER AS $$ BEGIN NEW.updated_at = CURRENT_TIMESTAMP; RETURN NEW; END; $$ LANGUAGE plpgsql; DROP TRIGGER IF EXISTS trigger_touch_outbox_updated_at ON public.outbox_events; CREATE TRIGGER trigger_touch_outbox_updated_at BEFORE UPDATE ON public.outbox_events FOR EACH ROW EXECUTE FUNCTION public.touch_outbox_updated_at(); """ ) def downgrade() -> None: # Event history remains available if the application is rolled back. pass