outbox.py 5.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184
  1. from __future__ import annotations
  2. from datetime import UTC, datetime
  3. from typing import Any
  4. from sqlalchemy import text
  5. from app.core.common.identifiers import new_governance_uid
  6. from app.core.events.consumer import ClaimedEvent, DispatchResult
  7. def enqueue_outbox(
  8. session: Any,
  9. *,
  10. aggregate_type: str,
  11. aggregate_id: str,
  12. event_type: str,
  13. payload: dict[str, Any],
  14. correlation_id: str | None = None,
  15. event_id: str | None = None,
  16. ) -> str:
  17. event_id = event_id or new_governance_uid()
  18. session.execute(
  19. text(
  20. """
  21. INSERT INTO public.outbox_events (
  22. event_id, aggregate_type, aggregate_id, event_type,
  23. payload, correlation_id
  24. ) VALUES (
  25. CAST(:event_id AS uuid), :aggregate_type, :aggregate_id,
  26. :event_type, CAST(:payload AS jsonb), :correlation_id
  27. )
  28. """
  29. ),
  30. {
  31. "event_id": event_id,
  32. "aggregate_type": aggregate_type,
  33. "aggregate_id": aggregate_id,
  34. "event_type": event_type,
  35. "payload": __import__("json").dumps(payload, ensure_ascii=False),
  36. "correlation_id": correlation_id,
  37. },
  38. )
  39. return event_id
  40. def claim_outbox(
  41. session: Any,
  42. limit: int = 50,
  43. event_types: list[str] | None = None,
  44. ) -> list[ClaimedEvent]:
  45. if event_types is not None and not event_types:
  46. return []
  47. rows = session.execute(
  48. text(
  49. """
  50. WITH selected AS (
  51. SELECT event_id
  52. FROM public.outbox_events
  53. WHERE status = 'pending' AND available_at <= CURRENT_TIMESTAMP
  54. AND (
  55. CAST(:event_types AS text[]) IS NULL
  56. OR event_type = ANY(CAST(:event_types AS text[]))
  57. )
  58. ORDER BY available_at, created_at
  59. FOR UPDATE SKIP LOCKED
  60. LIMIT :limit
  61. )
  62. UPDATE public.outbox_events AS event
  63. SET status = 'processing'
  64. FROM selected
  65. WHERE event.event_id = selected.event_id
  66. RETURNING event.event_id::text, event.event_type,
  67. event.payload, event.attempts
  68. """
  69. ),
  70. {
  71. "limit": int(limit),
  72. "event_types": event_types,
  73. },
  74. )
  75. return [
  76. ClaimedEvent(
  77. event_id=row[0],
  78. event_type=row[1],
  79. payload=dict(row[2]),
  80. attempts=int(row[3]),
  81. )
  82. for row in rows
  83. ]
  84. def consumed(session: Any, consumer_name: str, event_id: str) -> bool:
  85. return bool(
  86. session.execute(
  87. text(
  88. "SELECT 1 FROM public.outbox_consumptions "
  89. "WHERE consumer_name = :consumer_name "
  90. "AND event_id = CAST(:event_id AS uuid)"
  91. ),
  92. {"consumer_name": consumer_name, "event_id": event_id},
  93. ).scalar()
  94. )
  95. def apply_dispatch_result(
  96. session: Any,
  97. *,
  98. consumer_name: str,
  99. event: ClaimedEvent,
  100. result: DispatchResult,
  101. ) -> None:
  102. if result.record_consumption:
  103. session.execute(
  104. text(
  105. """
  106. INSERT INTO public.outbox_consumptions (consumer_name, event_id)
  107. VALUES (:consumer_name, CAST(:event_id AS uuid))
  108. ON CONFLICT (consumer_name, event_id) DO NOTHING
  109. """
  110. ),
  111. {"consumer_name": consumer_name, "event_id": event.event_id},
  112. )
  113. session.execute(
  114. text(
  115. """
  116. UPDATE public.outbox_events
  117. SET status = :status,
  118. attempts = :attempts,
  119. available_at = COALESCE(:available_at, available_at),
  120. published_at = CASE WHEN :status = 'published'
  121. THEN CURRENT_TIMESTAMP ELSE published_at END,
  122. last_error = :last_error
  123. WHERE event_id = CAST(:event_id AS uuid)
  124. """
  125. ),
  126. {
  127. "status": result.status,
  128. "attempts": result.attempts,
  129. "available_at": result.available_at,
  130. "last_error": result.last_error,
  131. "event_id": event.event_id,
  132. },
  133. )
  134. def outbox_health(session: Any) -> dict[str, Any]:
  135. row = session.execute(
  136. text(
  137. """
  138. SELECT
  139. COUNT(*) FILTER (WHERE status IN ('pending', 'processing')) AS backlog,
  140. COUNT(*) FILTER (WHERE status IN ('failed', 'dead_letter')) AS failed,
  141. EXTRACT(EPOCH FROM (
  142. CURRENT_TIMESTAMP - MIN(created_at)
  143. FILTER (WHERE status IN ('pending', 'processing'))
  144. )) AS oldest_age_seconds
  145. FROM public.outbox_events
  146. """
  147. )
  148. ).one()
  149. return {
  150. "backlog": int(row[0] or 0),
  151. "failed": int(row[1] or 0),
  152. "oldest_age_seconds": float(row[2] or 0),
  153. "checked_at": datetime.now(UTC).isoformat(),
  154. }
  155. def reset_stale_processing(session: Any, *, older_than_seconds: int = 300) -> int:
  156. result = session.execute(
  157. text(
  158. """
  159. UPDATE public.outbox_events
  160. SET status = 'pending', available_at = CURRENT_TIMESTAMP
  161. WHERE status = 'processing'
  162. AND updated_at < CURRENT_TIMESTAMP
  163. - (:older_than_seconds * INTERVAL '1 second')
  164. """
  165. ),
  166. {"older_than_seconds": int(older_than_seconds)},
  167. )
  168. return int(result.rowcount or 0)