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 # The token is stored inside the internal payload (the wire envelope below # explicitly selects only contract fields). It lets the outcome transaction # prove that this is still the same lease after network I/O. A stale worker # must never overwrite a later reclaim/retry or an idempotent callback. row.payload_json = { **row.payload_json, "_delivery_claim_token": str(uuid.uuid4()), } 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 claim_token = event.payload_json.get("_delivery_claim_token") if event.delivery_status != "delivering" or not isinstance(claim_token, str): return 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): acknowledged = body.get("ok") is True response_event_id = body.get("event_id") event_id_matches = response_event_id == str(event_id) result = body.get("result") execution_id = result.get("execution_id") if isinstance(result, dict) else None execution_id_valid = isinstance(execution_id, str) and bool(execution_id.strip()) success = acknowledged and event_id_matches and execution_id_valid if success: error = None error_code = None elif not acknowledged: error = ( "n8n response did not explicitly acknowledge the event with ok=true: " f"{body}" ) error_code = ( "remoteReportedFailure" if body.get("ok") is False else "malformedResponse" ) elif not event_id_matches: error = ( "n8n acknowledged a different event ID " f"(expected {event_id}, received {response_event_id!r})" ) error_code = "mismatchedEventId" else: error = "n8n response omitted a valid result.execution_id" error_code = "malformedResponse" 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 = ( f"Unexpected non-JSON-object response from n8n (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.scalar( select(OutboxEvent).where(OutboxEvent.event_id == event_id).with_for_update() ) if event is None: return if ( event.delivery_status != "delivering" or event.attempts != attempts or event.payload_json.get("_delivery_claim_token") != claim_token ): logger.info( "Ignoring stale delivery outcome for event %s because lease ownership changed", event_id, ) return event.attempts = attempts + 1 event.payload_json = { key: value for key, value in event.payload_json.items() if key != "_delivery_claim_token" } if success: event.delivery_status = "succeeded" event.last_error = None event.last_error_code = None event.next_attempt_at = None result = (body or {}).get("result") execution_id = result.get("execution_id") if isinstance(result, dict) else None event.external_run_id = execution_id if isinstance(execution_id, str) else None 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()