diff --git a/backend/app/core/config.py b/backend/app/core/config.py index 02cf4c3..14470d6 100644 --- a/backend/app/core/config.py +++ b/backend/app/core/config.py @@ -22,6 +22,7 @@ class Settings(BaseSettings): n8n_dispatch_interval_seconds: float = 3.0 n8n_http_timeout_seconds: float = 5.0 n8n_max_attempts: int = 5 + n8n_delivery_lease_seconds: float = 120.0 app_secret: str = "replace-in-production" session_cookie_name: str = "mobilityops_session" session_ttl_seconds: int = 60 * 60 * 8 diff --git a/backend/app/services/dispatcher.py b/backend/app/services/dispatcher.py index 67d59a5..c401002 100644 --- a/backend/app/services/dispatcher.py +++ b/backend/app/services/dispatcher.py @@ -22,8 +22,42 @@ 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).""" + """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) @@ -38,8 +72,10 @@ def _claim_due_events(batch_size: int = 5) -> list[uuid.UUID]: .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: @@ -120,7 +156,8 @@ def _deliver_one(event_id: uuid.UUID) -> None: def run_dispatch_cycle() -> int: - """Run one claim+deliver cycle. Returns the number of events processed.""" + """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) diff --git a/backend/tests/test_dispatcher.py b/backend/tests/test_dispatcher.py index 95e00cb..efab55d 100644 --- a/backend/tests/test_dispatcher.py +++ b/backend/tests/test_dispatcher.py @@ -1,7 +1,7 @@ from __future__ import annotations import uuid -from datetime import UTC, datetime +from datetime import UTC, datetime, timedelta from types import SimpleNamespace from sqlalchemy import select @@ -161,6 +161,83 @@ def test_deliver_one_handles_malformed_payload_without_getting_stuck(monkeypatch assert "Malformed outbox payload" in event.last_error +def test_claim_sets_a_lease_deadline(): + event_id = _make_pending_event("MO-006") + settings = get_settings() + before = datetime.now(UTC) + dispatcher._claim_due_events() + + event = _get_event(event_id) + assert event.delivery_status == "delivering" + assert event.next_attempt_at is not None + lease = settings.n8n_delivery_lease_seconds + assert event.next_attempt_at > before + timedelta(seconds=lease - 5) + + +def test_reclaim_ignores_an_active_unexpired_lease(): + # A worker that is still within its lease window must not be disturbed -- this is + # what prevents double delivery of an event another (still-alive) worker is handling. + event_id = _make_pending_event("MO-007") + dispatcher._claim_due_events() + + reclaimed = dispatcher._reclaim_stale_deliveries() + assert reclaimed == 0 + assert _get_event(event_id).delivery_status == "delivering" + + +def test_reclaim_recovers_an_expired_lease_and_preserves_attempts(monkeypatch): + # Simulates a process crash: the row was claimed (delivering) but no outcome was ever + # recorded, and its lease has since expired. + event_id = _make_pending_event("MO-008") + dispatcher._claim_due_events() + + db = SessionLocal() + try: + event = db.scalar(select(OutboxEvent).where(OutboxEvent.event_id == event_id)) + event.attempts = 2 + event.next_attempt_at = datetime.now(UTC) - timedelta(seconds=1) + db.commit() + finally: + db.close() + + reclaimed = dispatcher._reclaim_stale_deliveries() + assert reclaimed == 1 + + event = _get_event(event_id) + assert event.delivery_status == "pending" + assert event.next_attempt_at is None + assert event.attempts == 2 + assert "stale" in event.last_error.lower() + + # The reclaimed event is now a normal pending event, immediately claimable again. + claimed = dispatcher._claim_due_events() + assert event_id in claimed + + +def test_run_dispatch_cycle_recovers_a_stale_lease_before_claiming(monkeypatch): + event_id = _make_pending_event("MO-009") + dispatcher._claim_due_events() + db = SessionLocal() + try: + event = db.scalar(select(OutboxEvent).where(OutboxEvent.event_id == event_id)) + event.next_attempt_at = datetime.now(UTC) - timedelta(seconds=1) + db.commit() + finally: + db.close() + + def fake_post(url, json, timeout): + return SimpleNamespace( + raise_for_status=lambda: None, + json=lambda: {"ok": True, "event_id": str(event_id), "result": {}}, + ) + + monkeypatch.setattr(dispatcher.httpx, "post", fake_post) + processed = dispatcher.run_dispatch_cycle() + + assert processed >= 1 + assert _get_event(event_id).delivery_status == "succeeded" + + def test_run_dispatch_cycle_end_to_end(monkeypatch): event_id = _make_pending_event("MO-005")