test_cross_store_consistency.py 6.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194
  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.execute(
  76. text(
  77. "UPDATE public.outbox_events "
  78. "SET available_at = TIMESTAMPTZ '1970-01-01' "
  79. "WHERE event_id = CAST(:event_id AS uuid)"
  80. ),
  81. {"event_id": event_id},
  82. )
  83. session.commit()
  84. with Session(engine) as session:
  85. event = claim_outbox(session, limit=1)[0]
  86. result = dispatch_event(
  87. event,
  88. lambda _payload: (_ for _ in ()).throw(
  89. ConnectionError("downstream unavailable")
  90. ),
  91. max_attempts=3,
  92. now=datetime.now(timezone.utc),
  93. )
  94. apply_dispatch_result(
  95. session,
  96. consumer_name="integration-consumer",
  97. event=event,
  98. result=result,
  99. )
  100. session.commit()
  101. with Session(engine) as session:
  102. session.execute(
  103. text(
  104. "UPDATE public.outbox_events "
  105. "SET available_at = TIMESTAMPTZ '1970-01-01' "
  106. "WHERE event_id = CAST(:event_id AS uuid)"
  107. ),
  108. {"event_id": event_id},
  109. )
  110. session.commit()
  111. # A fresh Session represents a consumer restart.
  112. with Session(engine) as session:
  113. event = claim_outbox(session, limit=1)[0]
  114. result = dispatch_event(event, lambda payload: delivered.append(payload))
  115. apply_dispatch_result(
  116. session,
  117. consumer_name="integration-consumer",
  118. event=event,
  119. result=result,
  120. )
  121. session.commit()
  122. with Session(engine) as session:
  123. session.execute(
  124. text(
  125. "UPDATE public.outbox_events "
  126. "SET status = 'pending', "
  127. "available_at = TIMESTAMPTZ '1970-01-01' "
  128. "WHERE event_id = CAST(:event_id AS uuid)"
  129. ),
  130. {"event_id": event_id},
  131. )
  132. session.commit()
  133. duplicate_calls = []
  134. with Session(engine) as session:
  135. event = claim_outbox(session, limit=1)[0]
  136. result = dispatch_event(
  137. event,
  138. lambda payload: duplicate_calls.append(payload),
  139. already_processed=consumed(
  140. session, "integration-consumer", event.event_id
  141. ),
  142. )
  143. apply_dispatch_result(
  144. session,
  145. consumer_name="integration-consumer",
  146. event=event,
  147. result=result,
  148. )
  149. session.commit()
  150. assert delivered == [{"marker": marker}]
  151. assert duplicate_calls == []
  152. with engine.connect() as connection:
  153. row = connection.execute(
  154. text(
  155. "SELECT status, attempts FROM public.outbox_events "
  156. "WHERE event_id = CAST(:event_id AS uuid)"
  157. ),
  158. {"event_id": event_id},
  159. ).one()
  160. assert row[0] == "published"
  161. assert row[1] == 2
  162. finally:
  163. with engine.begin() as connection:
  164. connection.execute(
  165. text(
  166. "DELETE FROM public.outbox_consumptions "
  167. "WHERE event_id IN (SELECT event_id FROM public.outbox_events "
  168. "WHERE aggregate_id = :marker)"
  169. ),
  170. {"marker": marker},
  171. )
  172. connection.execute(
  173. text("DELETE FROM public.outbox_events WHERE aggregate_id = :marker"),
  174. {"marker": marker},
  175. )
  176. engine.dispose()