Files
MobilityOps/backend/app/services/dispatcher.py
T
NuklearRabbit 81e3fd63bd
MobilityOps acceptance / backend (push) Failing after 19s
MobilityOps acceptance / frontend (push) Successful in 25s
MobilityOps acceptance / e2e (push) Skipped
M54: harden operations and demo resilience
2026-08-24 03:31:03 +02:00

263 lines
10 KiB
Python

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()