outbox.py 5.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171
  1. from __future__ import annotations
  2. from datetime import datetime, timezone
  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(session: Any, limit: int = 50) -> list[ClaimedEvent]:
  41. rows = session.execute(
  42. text(
  43. """
  44. WITH selected AS (
  45. SELECT event_id
  46. FROM public.outbox_events
  47. WHERE status = 'pending' AND available_at <= CURRENT_TIMESTAMP
  48. ORDER BY available_at, created_at
  49. FOR UPDATE SKIP LOCKED
  50. LIMIT :limit
  51. )
  52. UPDATE public.outbox_events AS event
  53. SET status = 'processing'
  54. FROM selected
  55. WHERE event.event_id = selected.event_id
  56. RETURNING event.event_id::text, event.event_type,
  57. event.payload, event.attempts
  58. """
  59. ),
  60. {"limit": int(limit)},
  61. )
  62. return [
  63. ClaimedEvent(
  64. event_id=row[0],
  65. event_type=row[1],
  66. payload=dict(row[2]),
  67. attempts=int(row[3]),
  68. )
  69. for row in rows
  70. ]
  71. def consumed(session: Any, consumer_name: str, event_id: str) -> bool:
  72. return bool(
  73. session.execute(
  74. text(
  75. "SELECT 1 FROM public.outbox_consumptions "
  76. "WHERE consumer_name = :consumer_name "
  77. "AND event_id = CAST(:event_id AS uuid)"
  78. ),
  79. {"consumer_name": consumer_name, "event_id": event_id},
  80. ).scalar()
  81. )
  82. def apply_dispatch_result(
  83. session: Any,
  84. *,
  85. consumer_name: str,
  86. event: ClaimedEvent,
  87. result: DispatchResult,
  88. ) -> None:
  89. if result.record_consumption:
  90. session.execute(
  91. text(
  92. """
  93. INSERT INTO public.outbox_consumptions (consumer_name, event_id)
  94. VALUES (:consumer_name, CAST(:event_id AS uuid))
  95. ON CONFLICT (consumer_name, event_id) DO NOTHING
  96. """
  97. ),
  98. {"consumer_name": consumer_name, "event_id": event.event_id},
  99. )
  100. session.execute(
  101. text(
  102. """
  103. UPDATE public.outbox_events
  104. SET status = :status,
  105. attempts = :attempts,
  106. available_at = COALESCE(:available_at, available_at),
  107. published_at = CASE WHEN :status = 'published'
  108. THEN CURRENT_TIMESTAMP ELSE published_at END,
  109. last_error = :last_error
  110. WHERE event_id = CAST(:event_id AS uuid)
  111. """
  112. ),
  113. {
  114. "status": result.status,
  115. "attempts": result.attempts,
  116. "available_at": result.available_at,
  117. "last_error": result.last_error,
  118. "event_id": event.event_id,
  119. },
  120. )
  121. def outbox_health(session: Any) -> dict[str, Any]:
  122. row = session.execute(
  123. text(
  124. """
  125. SELECT
  126. COUNT(*) FILTER (WHERE status IN ('pending', 'processing')) AS backlog,
  127. COUNT(*) FILTER (WHERE status IN ('failed', 'dead_letter')) AS failed,
  128. EXTRACT(EPOCH FROM (
  129. CURRENT_TIMESTAMP - MIN(created_at)
  130. FILTER (WHERE status IN ('pending', 'processing'))
  131. )) AS oldest_age_seconds
  132. FROM public.outbox_events
  133. """
  134. )
  135. ).one()
  136. return {
  137. "backlog": int(row[0] or 0),
  138. "failed": int(row[1] or 0),
  139. "oldest_age_seconds": float(row[2] or 0),
  140. "checked_at": datetime.now(timezone.utc).isoformat(),
  141. }
  142. def reset_stale_processing(session: Any, *, older_than_seconds: int = 300) -> int:
  143. result = session.execute(
  144. text(
  145. """
  146. UPDATE public.outbox_events
  147. SET status = 'pending', available_at = CURRENT_TIMESTAMP
  148. WHERE status = 'processing'
  149. AND updated_at < CURRENT_TIMESTAMP
  150. - (:older_than_seconds * INTERVAL '1 second')
  151. """
  152. ),
  153. {"older_than_seconds": int(older_than_seconds)},
  154. )
  155. return int(result.rowcount or 0)