2026-07-29 22:50:27 +02:00
""" 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
2026-07-31 19:16:11 +02:00
op . execute ( " ALTER TABLE IF EXISTS event_outbox ADD COLUMN IF NOT EXISTS aggregate_type VARCHAR(100) " )
op . execute ( " ALTER TABLE IF EXISTS event_outbox ADD COLUMN IF NOT EXISTS aggregate_id UUID " )
op . execute ( " ALTER TABLE IF EXISTS event_outbox ADD COLUMN IF NOT EXISTS occurred_at TIMESTAMPTZ NOT NULL DEFAULT NOW() " )
op . execute ( " ALTER TABLE IF EXISTS event_outbox ADD COLUMN IF NOT EXISTS correlation_id UUID " )
op . execute ( " ALTER TABLE IF EXISTS event_outbox ADD COLUMN IF NOT EXISTS schema_version INTEGER NOT NULL DEFAULT 1 " )
2026-07-29 22:50:27 +02:00
2026-07-31 19:16:11 +02:00
op . execute ( " DO $$ BEGIN IF EXISTS (SELECT 1 FROM information_schema.tables WHERE table_name = ' event_outbox ' ) THEN CREATE INDEX IF NOT EXISTS ix_event_outbox_aggregate ON event_outbox (tenant_id, aggregate_type, aggregate_id); END IF; END $$ " )
op . execute ( " DO $$ BEGIN IF EXISTS (SELECT 1 FROM information_schema.tables WHERE table_name = ' event_outbox ' ) THEN CREATE INDEX IF NOT EXISTS ix_event_outbox_correlation ON event_outbox (correlation_id); END IF; END $$ " )
2026-07-29 22:50:27 +02:00
# 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 " ) ,
)
2026-07-31 19:16:11 +02:00
op . execute ( ' CREATE INDEX IF NOT EXISTS ix_outbox_deliveries_event ON outbox_deliveries (event_id) ' )
op . execute ( ' CREATE INDEX IF NOT EXISTS ix_outbox_deliveries_status ON outbox_deliveries (status, next_attempt_at) ' )
2026-07-29 22:50:27 +02:00
# 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 " )