| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185 |
- 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()
|