test_cross_store_consistency.py 6.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185
  1. from __future__ import annotations
  2. import os
  3. import uuid
  4. from datetime import datetime, timezone
  5. import pytest
  6. from sqlalchemy import create_engine, text
  7. from sqlalchemy.orm import Session
  8. from app.core.events.consumer import dispatch_event
  9. from app.core.events.outbox import (
  10. apply_dispatch_result,
  11. claim_outbox,
  12. consumed,
  13. enqueue_outbox,
  14. )
  15. pytestmark = pytest.mark.integration
  16. @pytest.fixture()
  17. def database_url():
  18. value = os.environ.get("TEST_DATABASE_URL")
  19. if not value:
  20. pytest.skip("TEST_DATABASE_URL is not configured")
  21. return value
  22. def test_command_and_event_rollback_atomically(database_url):
  23. engine = create_engine(database_url)
  24. marker = uuid.uuid4().hex
  25. event_id = None
  26. try:
  27. with Session(engine) as session:
  28. with pytest.raises(RuntimeError, match="force rollback"):
  29. with session.begin():
  30. session.execute(
  31. text(
  32. "INSERT INTO public.task_list "
  33. "(task_name, task_description, create_by) "
  34. "VALUES (:name, 'atomic test', 'pytest')"
  35. ),
  36. {"name": marker},
  37. )
  38. event_id = enqueue_outbox(
  39. session,
  40. aggregate_type="test_task",
  41. aggregate_id=marker,
  42. event_type="test.created",
  43. payload={"marker": marker},
  44. )
  45. raise RuntimeError("force rollback")
  46. with engine.connect() as connection:
  47. task_count = connection.execute(
  48. text("SELECT COUNT(*) FROM public.task_list WHERE task_name = :name"),
  49. {"name": marker},
  50. ).scalar_one()
  51. event_count = connection.execute(
  52. text(
  53. "SELECT COUNT(*) FROM public.outbox_events "
  54. "WHERE event_id = CAST(:event_id AS uuid)"
  55. ),
  56. {"event_id": event_id},
  57. ).scalar_one()
  58. assert task_count == 0
  59. assert event_count == 0
  60. finally:
  61. engine.dispose()
  62. def test_downstream_failure_recovers_and_restart_does_not_duplicate(database_url):
  63. engine = create_engine(database_url)
  64. marker = uuid.uuid4().hex
  65. delivered = []
  66. try:
  67. with Session(engine) as session:
  68. event_id = enqueue_outbox(
  69. session,
  70. aggregate_type="governance_object",
  71. aggregate_id=marker,
  72. event_type="governance.updated",
  73. payload={"marker": marker},
  74. )
  75. session.commit()
  76. with Session(engine) as session:
  77. event = claim_outbox(session, limit=1)[0]
  78. result = dispatch_event(
  79. event,
  80. lambda _payload: (_ for _ in ()).throw(
  81. ConnectionError("downstream unavailable")
  82. ),
  83. max_attempts=3,
  84. now=datetime.now(timezone.utc),
  85. )
  86. apply_dispatch_result(
  87. session,
  88. consumer_name="integration-consumer",
  89. event=event,
  90. result=result,
  91. )
  92. session.commit()
  93. with Session(engine) as session:
  94. session.execute(
  95. text(
  96. "UPDATE public.outbox_events "
  97. "SET available_at = CURRENT_TIMESTAMP "
  98. "WHERE event_id = CAST(:event_id AS uuid)"
  99. ),
  100. {"event_id": event_id},
  101. )
  102. session.commit()
  103. # A fresh Session represents a consumer restart.
  104. with Session(engine) as session:
  105. event = claim_outbox(session, limit=1)[0]
  106. result = dispatch_event(event, lambda payload: delivered.append(payload))
  107. apply_dispatch_result(
  108. session,
  109. consumer_name="integration-consumer",
  110. event=event,
  111. result=result,
  112. )
  113. session.commit()
  114. with Session(engine) as session:
  115. session.execute(
  116. text(
  117. "UPDATE public.outbox_events "
  118. "SET status = 'pending', available_at = CURRENT_TIMESTAMP "
  119. "WHERE event_id = CAST(:event_id AS uuid)"
  120. ),
  121. {"event_id": event_id},
  122. )
  123. session.commit()
  124. duplicate_calls = []
  125. with Session(engine) as session:
  126. event = claim_outbox(session, limit=1)[0]
  127. result = dispatch_event(
  128. event,
  129. lambda payload: duplicate_calls.append(payload),
  130. already_processed=consumed(
  131. session, "integration-consumer", event.event_id
  132. ),
  133. )
  134. apply_dispatch_result(
  135. session,
  136. consumer_name="integration-consumer",
  137. event=event,
  138. result=result,
  139. )
  140. session.commit()
  141. assert delivered == [{"marker": marker}]
  142. assert duplicate_calls == []
  143. with engine.connect() as connection:
  144. row = connection.execute(
  145. text(
  146. "SELECT status, attempts FROM public.outbox_events "
  147. "WHERE event_id = CAST(:event_id AS uuid)"
  148. ),
  149. {"event_id": event_id},
  150. ).one()
  151. assert row[0] == "published"
  152. assert row[1] == 2
  153. finally:
  154. with engine.begin() as connection:
  155. connection.execute(
  156. text(
  157. "DELETE FROM public.outbox_consumptions "
  158. "WHERE event_id IN (SELECT event_id FROM public.outbox_events "
  159. "WHERE aggregate_id = :marker)"
  160. ),
  161. {"marker": marker},
  162. )
  163. connection.execute(
  164. text("DELETE FROM public.outbox_events WHERE aggregate_id = :marker"),
  165. {"marker": marker},
  166. )
  167. engine.dispose()