from __future__ import annotations import logging import threading import uuid from datetime import UTC, datetime, timedelta import httpx from sqlalchemy import select from app.core.config import get_settings from app.core.db import SessionLocal from app.models.outbox import OutboxEvent logger = logging.getLogger("mobilityops.dispatcher") settings = get_settings() _stop_event = threading.Event() def _backoff_seconds(attempts: int) -> int: return min(2**attempts, 60) def _reclaim_stale_deliveries(batch_size: int = 10) -> int: """Recover events stuck in 'delivering' because the process that claimed them died before recording an outcome. Only leases whose deadline has passed are touched, so an in-flight delivery from a still-alive worker is never disturbed or double-processed; `attempts` is preserved so the count reflects true history.""" db = SessionLocal() try: now = datetime.now(UTC) rows = db.scalars( select(OutboxEvent) .where( OutboxEvent.delivery_status == "delivering", OutboxEvent.next_attempt_at.is_not(None), OutboxEvent.next_attempt_at <= now, ) .limit(batch_size) .with_for_update(skip_locked=True) ).all() for row in rows: row.delivery_status = "pending" row.next_attempt_at = None row.last_error = ( "Recovered from a stale 'delivering' lease " f"(no outcome recorded within {settings.n8n_delivery_lease_seconds:.0f}s; " f"the process likely crashed mid-delivery). attempts preserved at {row.attempts}." )[:2000] row.last_error_code = "staleLeaseRecovered" db.commit() return len(rows) finally: db.close() def _claim_due_events(batch_size: int = 5) -> list[uuid.UUID]: """Claim a batch of due events with a short-lived transaction (no network I/O held open). Each claimed row gets a lease deadline (next_attempt_at) so a crash between this claim and the outcome being recorded is recoverable by _reclaim_stale_deliveries.""" db = SessionLocal() try: now = datetime.now(UTC) rows = db.scalars( select(OutboxEvent) .where( OutboxEvent.delivery_status == "pending", (OutboxEvent.next_attempt_at.is_(None)) | (OutboxEvent.next_attempt_at <= now), ) .order_by(OutboxEvent.occurred_at.asc()) .limit(batch_size) .with_for_update(skip_locked=True) ).all() claimed_ids = [row.event_id for row in rows] lease_deadline = now + timedelta(seconds=settings.n8n_delivery_lease_seconds) for row in rows: row.delivery_status = "delivering" row.next_attempt_at = lease_deadline db.commit() return claimed_ids finally: db.close() def _deliver_one(event_id: uuid.UUID) -> None: db = SessionLocal() try: event = db.get(OutboxEvent, event_id) if event is None: return # Reconstruct the wire envelope from contracts/events.schema.json: only the fields # the schema declares (additionalProperties: false), sourced from real columns where # possible. `payload_json` also carries an internal `aggregate_ref` convenience field # for our own dashboard/audit reads, which must not be forwarded to n8n. try: wire_event = { "event_id": str(event.event_id), "event_type": event.event_type, "occurred_at": event.occurred_at.isoformat(), "correlation_id": event.payload_json["correlation_id"], "aggregate": event.payload_json["aggregate"], "data": event.payload_json["data"], } payload_error: str | None = None except KeyError as exc: # A malformed payload must still resolve the claimed "delivering" row to a # terminal-or-retryable state below, rather than leaving it stuck forever. wire_event = None payload_error = f"Malformed outbox payload, missing key {exc}" attempts = event.attempts finally: db.close() error_code: str | None if wire_event is None: success, error, body = False, payload_error, None error_code = "malformedPayload" else: try: response = httpx.post( settings.n8n_webhook_url, json=wire_event, headers={"X-Fleet-Ops-Trigger-Token": settings.n8n_webhook_trigger_token}, timeout=settings.n8n_http_timeout_seconds, ) response.raise_for_status() try: body = response.json() except ValueError: body = None if isinstance(body, dict): success = bool(body.get("ok", True)) error = None if success else f"n8n reported failure: {body}" error_code = None if success else "remoteReportedFailure" else: # A 2xx status with a non-object (or unparsable) body means the workflow # itself errored before its "Respond to Webhook" node ran -- n8n's default # error response still carries a 2xx-looking status here. Treat it as a # failure so the event is retried rather than lost or wrongly marked # succeeded. success = False error = ( "Unexpected non-JSON-object response from n8n " f"(status {response.status_code})" ) error_code = "malformedResponse" except httpx.HTTPError as exc: success = False error = f"{type(exc).__name__}: {exc}" error_code = "connectionError" body = None db = SessionLocal() try: event = db.get(OutboxEvent, event_id) if event is None: return event.attempts = attempts + 1 if success: event.delivery_status = "succeeded" event.last_error = None event.last_error_code = None event.next_attempt_at = None event.external_run_id = str((body or {}).get("event_id", event_id)) else: event.last_error = (error or "delivery failed")[:2000] event.last_error_code = error_code or "unknownError" if event.attempts >= settings.n8n_max_attempts: event.delivery_status = "failed" event.next_attempt_at = None else: event.delivery_status = "pending" event.next_attempt_at = datetime.now(UTC) + timedelta( seconds=_backoff_seconds(event.attempts) ) db.commit() finally: db.close() def run_dispatch_cycle() -> int: """Run one reclaim+claim+deliver cycle. Returns the number of events processed.""" _reclaim_stale_deliveries() claimed = _claim_due_events() for event_id in claimed: _deliver_one(event_id) return len(claimed) def _loop() -> None: while not _stop_event.is_set(): try: run_dispatch_cycle() except Exception: # noqa: BLE001 logger.exception("Outbox dispatch cycle failed") _stop_event.wait(settings.n8n_dispatch_interval_seconds) def start_background_dispatcher() -> None: if not settings.n8n_dispatch_enabled: return _stop_event.clear() thread = threading.Thread(target=_loop, name="outbox-dispatcher", daemon=True) thread.start() def stop_background_dispatcher() -> None: _stop_event.set()