_claim_due_events flipped rows to 'delivering' and committed before the HTTP call; if the process died between that commit and the outcome- recording transaction, the row stayed 'delivering' forever with no reclaim path -- a real gap, not previously documented as an accepted limitation. Give each claim a lease deadline (reusing next_attempt_at, since it's only otherwise meaningful for pending-status backoff scheduling) and sweep expired leases back to pending at the start of every dispatch cycle, before claiming new work. attempts is preserved so the count still reflects true history. Only leases past their deadline are touched, so a still-alive worker mid-delivery is never disturbed or double-processed.
186 lines
6.5 KiB
Python
186 lines
6.5 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]
|
|
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()
|
|
|
|
if wire_event is None:
|
|
success, error, body = False, payload_error, None
|
|
else:
|
|
try:
|
|
response = httpx.post(
|
|
settings.n8n_webhook_url,
|
|
json=wire_event,
|
|
timeout=settings.n8n_http_timeout_seconds,
|
|
)
|
|
response.raise_for_status()
|
|
body = response.json()
|
|
success = bool(body.get("ok", True))
|
|
error = None if success else f"n8n reported failure: {body}"
|
|
except httpx.HTTPError as exc:
|
|
success = False
|
|
error = f"{type(exc).__name__}: {exc}"
|
|
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.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]
|
|
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()
|