| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171 |
- from __future__ import annotations
- from datetime import datetime, timezone
- from typing import Any
- from sqlalchemy import text
- from app.core.common.identifiers import new_governance_uid
- from app.core.events.consumer import ClaimedEvent, DispatchResult
- def enqueue_outbox(
- session: Any,
- *,
- aggregate_type: str,
- aggregate_id: str,
- event_type: str,
- payload: dict[str, Any],
- correlation_id: str | None = None,
- event_id: str | None = None,
- ) -> str:
- event_id = event_id or new_governance_uid()
- session.execute(
- text(
- """
- INSERT INTO public.outbox_events (
- event_id, aggregate_type, aggregate_id, event_type,
- payload, correlation_id
- ) VALUES (
- CAST(:event_id AS uuid), :aggregate_type, :aggregate_id,
- :event_type, CAST(:payload AS jsonb), :correlation_id
- )
- """
- ),
- {
- "event_id": event_id,
- "aggregate_type": aggregate_type,
- "aggregate_id": aggregate_id,
- "event_type": event_type,
- "payload": __import__("json").dumps(payload, ensure_ascii=False),
- "correlation_id": correlation_id,
- },
- )
- return event_id
- def claim_outbox(session: Any, limit: int = 50) -> list[ClaimedEvent]:
- rows = session.execute(
- text(
- """
- WITH selected AS (
- SELECT event_id
- FROM public.outbox_events
- WHERE status = 'pending' AND available_at <= CURRENT_TIMESTAMP
- ORDER BY available_at, created_at
- FOR UPDATE SKIP LOCKED
- LIMIT :limit
- )
- UPDATE public.outbox_events AS event
- SET status = 'processing'
- FROM selected
- WHERE event.event_id = selected.event_id
- RETURNING event.event_id::text, event.event_type,
- event.payload, event.attempts
- """
- ),
- {"limit": int(limit)},
- )
- return [
- ClaimedEvent(
- event_id=row[0],
- event_type=row[1],
- payload=dict(row[2]),
- attempts=int(row[3]),
- )
- for row in rows
- ]
- def consumed(session: Any, consumer_name: str, event_id: str) -> bool:
- return bool(
- session.execute(
- text(
- "SELECT 1 FROM public.outbox_consumptions "
- "WHERE consumer_name = :consumer_name "
- "AND event_id = CAST(:event_id AS uuid)"
- ),
- {"consumer_name": consumer_name, "event_id": event_id},
- ).scalar()
- )
- def apply_dispatch_result(
- session: Any,
- *,
- consumer_name: str,
- event: ClaimedEvent,
- result: DispatchResult,
- ) -> None:
- if result.record_consumption:
- session.execute(
- text(
- """
- INSERT INTO public.outbox_consumptions (consumer_name, event_id)
- VALUES (:consumer_name, CAST(:event_id AS uuid))
- ON CONFLICT (consumer_name, event_id) DO NOTHING
- """
- ),
- {"consumer_name": consumer_name, "event_id": event.event_id},
- )
- session.execute(
- text(
- """
- UPDATE public.outbox_events
- SET status = :status,
- attempts = :attempts,
- available_at = COALESCE(:available_at, available_at),
- published_at = CASE WHEN :status = 'published'
- THEN CURRENT_TIMESTAMP ELSE published_at END,
- last_error = :last_error
- WHERE event_id = CAST(:event_id AS uuid)
- """
- ),
- {
- "status": result.status,
- "attempts": result.attempts,
- "available_at": result.available_at,
- "last_error": result.last_error,
- "event_id": event.event_id,
- },
- )
- def outbox_health(session: Any) -> dict[str, Any]:
- row = session.execute(
- text(
- """
- SELECT
- COUNT(*) FILTER (WHERE status IN ('pending', 'processing')) AS backlog,
- COUNT(*) FILTER (WHERE status IN ('failed', 'dead_letter')) AS failed,
- EXTRACT(EPOCH FROM (
- CURRENT_TIMESTAMP - MIN(created_at)
- FILTER (WHERE status IN ('pending', 'processing'))
- )) AS oldest_age_seconds
- FROM public.outbox_events
- """
- )
- ).one()
- return {
- "backlog": int(row[0] or 0),
- "failed": int(row[1] or 0),
- "oldest_age_seconds": float(row[2] or 0),
- "checked_at": datetime.now(timezone.utc).isoformat(),
- }
- def reset_stale_processing(session: Any, *, older_than_seconds: int = 300) -> int:
- result = session.execute(
- text(
- """
- UPDATE public.outbox_events
- SET status = 'pending', available_at = CURRENT_TIMESTAMP
- WHERE status = 'processing'
- AND updated_at < CURRENT_TIMESTAMP
- - (:older_than_seconds * INTERVAL '1 second')
- """
- ),
- {"older_than_seconds": int(older_than_seconds)},
- )
- return int(result.rowcount or 0)
|