Files
group_fqcd_jr/tests/integration/test_outbox_mysql.py

72 lines
2.4 KiB
Python

from datetime import UTC, datetime
from uuid import uuid4
import pytest
from sqlalchemy import delete, select
from app.infrastructure.db import SessionFactory
from app.model.platform import DomainEventOutbox, OutboxDelivery
from app.worker.outbox_worker import OutboxWorker
@pytest.mark.asyncio
async def test_duplicate_delivery_skips_handler_after_restart() -> None:
event_id = str(uuid4())
event_type = f"test.mysql.{uuid4().hex}"
now = datetime.now(UTC).replace(tzinfo=None)
calls = 0
async def handler(payload: dict[str, object]) -> None:
nonlocal calls
calls += 1
try:
async with SessionFactory() as session:
session.add(
DomainEventOutbox(
event_id=event_id,
event_type=event_type,
aggregate_type="test",
aggregate_id=event_id,
trace_id=event_id,
payload={"value": 1},
status="pending",
retry_count=0,
occurred_at=now,
created_at=now,
updated_at=now,
)
)
await session.commit()
async with SessionFactory() as session:
assert await OutboxWorker(session, {event_type: handler}).publish_one()
async with SessionFactory() as session:
event = await session.scalar(
select(DomainEventOutbox).where(DomainEventOutbox.event_id == event_id)
)
assert event is not None
event.status = "pending"
event.published_at = None
await session.commit()
async with SessionFactory() as session:
assert await OutboxWorker(session, {event_type: handler}).publish_one()
assert calls == 1
async with SessionFactory() as session:
deliveries = list(
await session.scalars(
select(OutboxDelivery).where(OutboxDelivery.event_id == event_id)
)
)
assert len(deliveries) == 1
finally:
async with SessionFactory() as session:
await session.execute(delete(OutboxDelivery).where(OutboxDelivery.event_id == event_id))
await session.execute(
delete(DomainEventOutbox).where(DomainEventOutbox.event_id == event_id)
)
await session.commit()