20260716_03_outbox.py 2.3 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465
  1. """Create the transactional outbox and idempotent consumption ledger."""
  2. from alembic import op
  3. revision = "20260716_03"
  4. down_revision = "20260716_02"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.execute(
  9. """
  10. CREATE TABLE IF NOT EXISTS public.outbox_events (
  11. event_id UUID PRIMARY KEY,
  12. aggregate_type VARCHAR(100) NOT NULL,
  13. aggregate_id VARCHAR(200) NOT NULL,
  14. event_type VARCHAR(150) NOT NULL,
  15. payload JSONB NOT NULL,
  16. correlation_id VARCHAR(100),
  17. status VARCHAR(30) NOT NULL DEFAULT 'pending',
  18. attempts INTEGER NOT NULL DEFAULT 0,
  19. available_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  20. created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  21. updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  22. published_at TIMESTAMPTZ,
  23. last_error TEXT,
  24. CONSTRAINT ck_outbox_status CHECK (
  25. status IN ('pending', 'processing', 'published',
  26. 'failed', 'dead_letter')
  27. )
  28. );
  29. CREATE INDEX IF NOT EXISTS idx_outbox_delivery
  30. ON public.outbox_events(status, available_at, created_at);
  31. CREATE INDEX IF NOT EXISTS idx_outbox_aggregate
  32. ON public.outbox_events(aggregate_type, aggregate_id);
  33. CREATE TABLE IF NOT EXISTS public.outbox_consumptions (
  34. consumer_name VARCHAR(100) NOT NULL,
  35. event_id UUID NOT NULL REFERENCES public.outbox_events(event_id),
  36. processed_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
  37. PRIMARY KEY (consumer_name, event_id)
  38. );
  39. CREATE OR REPLACE FUNCTION public.touch_outbox_updated_at()
  40. RETURNS TRIGGER AS $$
  41. BEGIN
  42. NEW.updated_at = CURRENT_TIMESTAMP;
  43. RETURN NEW;
  44. END;
  45. $$ LANGUAGE plpgsql;
  46. DROP TRIGGER IF EXISTS trigger_touch_outbox_updated_at
  47. ON public.outbox_events;
  48. CREATE TRIGGER trigger_touch_outbox_updated_at
  49. BEFORE UPDATE ON public.outbox_events
  50. FOR EACH ROW EXECUTE FUNCTION public.touch_outbox_updated_at();
  51. """
  52. )
  53. def downgrade() -> None:
  54. # Event history remains available if the application is rolled back.
  55. pass