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