From 0f4e51c4b3d69c68f2e5277c6f43b3f41cd11eda Mon Sep 17 00:00:00 2001 From: Agent Zero Date: Wed, 29 Jul 2026 12:49:15 +0200 Subject: [PATCH] fix: consumer_inbox table for outbox idempotency + tenant_plugin_activation table --- alembic/versions/0065_consumer_inbox.py | 40 +++++++++++++++++ .../versions/0066_tenant_plugin_activation.py | 44 +++++++++++++++++++ 2 files changed, 84 insertions(+) create mode 100644 alembic/versions/0065_consumer_inbox.py create mode 100644 alembic/versions/0066_tenant_plugin_activation.py diff --git a/alembic/versions/0065_consumer_inbox.py b/alembic/versions/0065_consumer_inbox.py new file mode 100644 index 0000000..f1734bc --- /dev/null +++ b/alembic/versions/0065_consumer_inbox.py @@ -0,0 +1,40 @@ +"""Add consumer_inbox table for outbox idempotency. + +Revision ID: 0065 +Revises: 0064 +Create Date: 2026-07-29 + +Without idempotency, a worker crash between sending an email/webhook +and marking the event as published can lead to duplicate deliveries. + +This migration creates a consumer_inbox table that tracks which +consumers have already processed which events. +""" + +from alembic import op +import sqlalchemy as sa +from sqlalchemy.dialects.postgresql import UUID + +revision = "0065" +down_revision = "0064" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.create_table( + "consumer_inbox", + sa.Column("id", UUID(as_uuid=True), primary_key=True, server_default=sa.text("gen_random_uuid()")), + sa.Column("event_id", UUID(as_uuid=True), sa.ForeignKey("event_outbox.id", ondelete="CASCADE"), nullable=False, index=True), + sa.Column("consumer_name", sa.String(100), nullable=False, index=True), + sa.Column("status", sa.String(20), nullable=False, default="pending"), # pending, processed, failed + sa.Column("processed_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("error_message", sa.Text, nullable=True), + sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.text("NOW()"), nullable=False), + sa.UniqueConstraint("event_id", "consumer_name", name="uq_consumer_inbox_event_consumer"), + ) + op.execute("ALTER TABLE consumer_inbox ENABLE ROW LEVEL SECURITY") + + +def downgrade() -> None: + op.drop_table("consumer_inbox") diff --git a/alembic/versions/0066_tenant_plugin_activation.py b/alembic/versions/0066_tenant_plugin_activation.py new file mode 100644 index 0000000..331ee65 --- /dev/null +++ b/alembic/versions/0066_tenant_plugin_activation.py @@ -0,0 +1,44 @@ +"""Add tenant_plugin_activation table for per-tenant plugin activation. + +Revision ID: 0066 +Revises: 0065 +Create Date: 2026-07-29 + +Currently plugins are activated globally. This migration creates a +table for per-tenant plugin activation so that different tenants can +enable/disable plugins independently. +""" + +from alembic import op +import sqlalchemy as sa +from sqlalchemy.dialects.postgresql import UUID + +revision = "0066" +down_revision = "0065" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.create_table( + "tenant_plugin_activation", + sa.Column("id", UUID(as_uuid=True), primary_key=True, server_default=sa.text("gen_random_uuid()")), + sa.Column("tenant_id", UUID(as_uuid=True), sa.ForeignKey("tenants.id", ondelete="CASCADE"), nullable=False, index=True), + sa.Column("plugin_name", sa.String(100), nullable=False, index=True), + sa.Column("is_active", sa.Boolean, nullable=False, default=True), + sa.Column("activated_by", UUID(as_uuid=True), sa.ForeignKey("users.id", ondelete="SET NULL"), 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("tenant_id", "plugin_name", name="uq_tenant_plugin"), + ) + op.execute("ALTER TABLE tenant_plugin_activation ENABLE ROW LEVEL SECURITY") + op.execute(""" + CREATE POLICY tenant_plugin_activation_tenant_isolation ON tenant_plugin_activation + FOR ALL + USING (tenant_id = current_setting('app.current_tenant_id', true)::uuid) + WITH CHECK (tenant_id = current_setting('app.current_tenant_id', true)::uuid) + """) + + +def downgrade() -> None: + op.drop_table("tenant_plugin_activation")