81 lines
4.0 KiB
Python
81 lines
4.0 KiB
Python
"""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.execute("ALTER TABLE event_outbox ADD COLUMN IF NOT EXISTS aggregate_type VARCHAR(100)")
|
|
op.execute("ALTER TABLE event_outbox ADD COLUMN IF NOT EXISTS aggregate_id UUID")
|
|
op.execute("ALTER TABLE event_outbox ADD COLUMN IF NOT EXISTS occurred_at TIMESTAMPTZ NOT NULL DEFAULT NOW()")
|
|
op.execute("ALTER TABLE event_outbox ADD COLUMN IF NOT EXISTS correlation_id UUID")
|
|
op.execute("ALTER TABLE event_outbox ADD COLUMN IF NOT EXISTS schema_version INTEGER NOT NULL DEFAULT 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")
|