from __future__ import annotations import os import uuid from datetime import datetime, timezone import pytest from sqlalchemy import create_engine, text from sqlalchemy.orm import Session from app.core.events.consumer import dispatch_event from app.core.events.outbox import ( apply_dispatch_result, claim_outbox, consumed, enqueue_outbox, ) pytestmark = pytest.mark.integration @pytest.fixture() def database_url(): value = os.environ.get("TEST_DATABASE_URL") if not value: pytest.skip("TEST_DATABASE_URL is not configured") return value def test_command_and_event_rollback_atomically(database_url): engine = create_engine(database_url) marker = uuid.uuid4().hex event_id = None try: with Session(engine) as session: with pytest.raises(RuntimeError, match="force rollback"): with session.begin(): session.execute( text( "INSERT INTO public.task_list " "(task_name, task_description, create_by) " "VALUES (:name, 'atomic test', 'pytest')" ), {"name": marker}, ) event_id = enqueue_outbox( session, aggregate_type="test_task", aggregate_id=marker, event_type="test.created", payload={"marker": marker}, ) raise RuntimeError("force rollback") with engine.connect() as connection: task_count = connection.execute( text("SELECT COUNT(*) FROM public.task_list WHERE task_name = :name"), {"name": marker}, ).scalar_one() event_count = connection.execute( text( "SELECT COUNT(*) FROM public.outbox_events " "WHERE event_id = CAST(:event_id AS uuid)" ), {"event_id": event_id}, ).scalar_one() assert task_count == 0 assert event_count == 0 finally: engine.dispose() def test_downstream_failure_recovers_and_restart_does_not_duplicate(database_url): engine = create_engine(database_url) marker = uuid.uuid4().hex delivered = [] try: with Session(engine) as session: event_id = enqueue_outbox( session, aggregate_type="governance_object", aggregate_id=marker, event_type="governance.updated", payload={"marker": marker}, ) session.commit() with Session(engine) as session: event = claim_outbox(session, limit=1)[0] result = dispatch_event( event, lambda _payload: (_ for _ in ()).throw( ConnectionError("downstream unavailable") ), max_attempts=3, now=datetime.now(timezone.utc), ) apply_dispatch_result( session, consumer_name="integration-consumer", event=event, result=result, ) session.commit() with Session(engine) as session: session.execute( text( "UPDATE public.outbox_events " "SET available_at = CURRENT_TIMESTAMP " "WHERE event_id = CAST(:event_id AS uuid)" ), {"event_id": event_id}, ) session.commit() # A fresh Session represents a consumer restart. with Session(engine) as session: event = claim_outbox(session, limit=1)[0] result = dispatch_event(event, lambda payload: delivered.append(payload)) apply_dispatch_result( session, consumer_name="integration-consumer", event=event, result=result, ) session.commit() with Session(engine) as session: session.execute( text( "UPDATE public.outbox_events " "SET status = 'pending', available_at = CURRENT_TIMESTAMP " "WHERE event_id = CAST(:event_id AS uuid)" ), {"event_id": event_id}, ) session.commit() duplicate_calls = [] with Session(engine) as session: event = claim_outbox(session, limit=1)[0] result = dispatch_event( event, lambda payload: duplicate_calls.append(payload), already_processed=consumed( session, "integration-consumer", event.event_id ), ) apply_dispatch_result( session, consumer_name="integration-consumer", event=event, result=result, ) session.commit() assert delivered == [{"marker": marker}] assert duplicate_calls == [] with engine.connect() as connection: row = connection.execute( text( "SELECT status, attempts FROM public.outbox_events " "WHERE event_id = CAST(:event_id AS uuid)" ), {"event_id": event_id}, ).one() assert row[0] == "published" assert row[1] == 2 finally: with engine.begin() as connection: connection.execute( text( "DELETE FROM public.outbox_consumptions " "WHERE event_id IN (SELECT event_id FROM public.outbox_events " "WHERE aggregate_id = :marker)" ), {"marker": marker}, ) connection.execute( text("DELETE FROM public.outbox_events WHERE aggregate_id = :marker"), {"marker": marker}, ) engine.dispose()