"""Add outbox_deliveries table and envelope columns to event_outbox. Standardized Event-Envelope: - event_id (already exists as id) - event_type (already exists as event_name) - tenant_id (already exists) - aggregate_type (NEW) - aggregate_id (NEW) - occurred_at (NEW) - correlation_id (NEW) - schema_version (NEW, default 1) - payload (already exists) outbox_deliveries tracks per-consumer delivery status. An event is only 'published' when all mandatory deliveries succeed. Revision ID: 0075 Revises: 0074 """ from alembic import op import sqlalchemy as sa from sqlalchemy.dialects.postgresql import UUID as PGUUID revision = "0075" down_revision = "0074" branch_labels = None depends_on = None def upgrade() -> None: # 1. Add envelope columns to event_outbox op.add_column("event_outbox", sa.Column("aggregate_type", sa.String(100), nullable=True)) op.add_column("event_outbox", sa.Column("aggregate_id", PGUUID(as_uuid=True), nullable=True)) op.add_column("event_outbox", sa.Column("occurred_at", sa.DateTime(timezone=True), server_default=sa.text("NOW()"), nullable=False)) op.add_column("event_outbox", sa.Column("correlation_id", PGUUID(as_uuid=True), nullable=True)) op.add_column("event_outbox", sa.Column("schema_version", sa.Integer, nullable=False, server_default=sa.text("1"))) op.execute("CREATE INDEX IF NOT EXISTS ix_event_outbox_aggregate ON event_outbox (tenant_id, aggregate_type, aggregate_id)") op.execute("CREATE INDEX IF NOT EXISTS ix_event_outbox_correlation ON event_outbox (correlation_id)") # 2. Create outbox_deliveries table op.create_table( "outbox_deliveries", sa.Column("id", PGUUID(as_uuid=True), primary_key=True, server_default=sa.text("gen_random_uuid()")), sa.Column("event_id", PGUUID(as_uuid=True), sa.ForeignKey("event_outbox.id", ondelete="CASCADE"), nullable=False), sa.Column("consumer_name", sa.String(150), nullable=False), sa.Column("status", sa.String(30), nullable=False, server_default="pending"), sa.Column("attempt_count", sa.Integer, nullable=False, server_default=sa.text("0")), sa.Column("next_attempt_at", sa.DateTime(timezone=True), nullable=True), sa.Column("last_error", sa.Text, nullable=True), sa.Column("processed_at", sa.DateTime(timezone=True), nullable=True), sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.text("NOW()"), nullable=False), sa.Column("updated_at", sa.DateTime(timezone=True), server_default=sa.text("NOW()"), nullable=False), sa.UniqueConstraint("event_id", "consumer_name", name="uq_outbox_deliveries_event_consumer"), ) op.create_index("ix_outbox_deliveries_event", "outbox_deliveries", ["event_id"]) op.create_index("ix_outbox_deliveries_status", "outbox_deliveries", ["status", "next_attempt_at"]) # RLS + Grants op.execute("ALTER TABLE outbox_deliveries ENABLE ROW LEVEL SECURITY") op.execute( "CREATE POLICY outbox_deliveries_tenant_isolation ON outbox_deliveries " "FOR ALL " "USING (EXISTS (SELECT 1 FROM event_outbox WHERE event_outbox.id = outbox_deliveries.event_id AND event_outbox.tenant_id = current_setting('app.current_tenant_id', true)::uuid)) " "WITH CHECK (EXISTS (SELECT 1 FROM event_outbox WHERE event_outbox.id = outbox_deliveries.event_id AND event_outbox.tenant_id = current_setting('app.current_tenant_id', true)::uuid))" ) op.execute("GRANT SELECT, INSERT, UPDATE, DELETE ON outbox_deliveries TO crm_api, crm_worker") def downgrade() -> None: op.execute("DROP POLICY IF EXISTS outbox_deliveries_tenant_isolation ON outbox_deliveries") op.drop_table("outbox_deliveries") op.execute("DROP INDEX IF EXISTS ix_event_outbox_correlation") op.execute("DROP INDEX IF EXISTS ix_event_outbox_aggregate") op.drop_column("event_outbox", "schema_version") op.drop_column("event_outbox", "correlation_id") op.drop_column("event_outbox", "occurred_at") op.drop_column("event_outbox", "aggregate_id") op.drop_column("event_outbox", "aggregate_type")