"""PostgreSQL-backed webhook delivery worker with bounded retries.""" from __future__ import annotations import asyncio from dataclasses import dataclass from datetime import UTC, datetime, timedelta from db.models import WebhookDelivery from db.session import SessionLocal from loguru import logger from services.webhooks.delivery import DeliveryResult, deliver_webhook, should_retry from services.webhooks.security import UnsafeWebhookUrl from sqlalchemy import or_, select, update POLL_INTERVAL_SECONDS = 2.0 MAX_ATTEMPTS = 5 RETRY_DELAYS_SECONDS = (30, 120, 600, 3600) @dataclass(frozen=True) class ClaimedDelivery: id: str event_id: str url: str secret: str payload: dict async def recover_interrupted_webhooks() -> None: async with SessionLocal() as session: await session.execute( update(WebhookDelivery) .where(WebhookDelivery.status == "processing") .values(status="pending", next_attempt_at=datetime.now(UTC)) ) await session.commit() async def claim_pending_webhook() -> ClaimedDelivery | None: now = datetime.now(UTC) async with SessionLocal() as session: row = ( await session.execute( select(WebhookDelivery) .where( WebhookDelivery.status == "pending", or_( WebhookDelivery.next_attempt_at.is_(None), WebhookDelivery.next_attempt_at <= now, ), ) .order_by(WebhookDelivery.next_attempt_at, WebhookDelivery.created_at) .with_for_update(skip_locked=True) .limit(1) ) ).scalar_one_or_none() if not row: return None row.status = "processing" await session.commit() return ClaimedDelivery( id=row.id, event_id=row.event_id, url=row.url, secret=str((row.secrets or {}).get("secret") or ""), payload=dict(row.payload or {}), ) async def finish_delivery(delivery_id: str, result: DeliveryResult) -> None: async with SessionLocal() as session: row = await session.get(WebhookDelivery, delivery_id) if not row or row.status != "processing": return row.attempt_count += 1 row.last_status_code = result.status_code row.last_error = result.error[:2048] if result.ok: row.status = "delivered" row.delivered_at = datetime.now(UTC) row.next_attempt_at = None elif should_retry(result.status_code) and row.attempt_count < MAX_ATTEMPTS: delay_index = min(row.attempt_count - 1, len(RETRY_DELAYS_SECONDS) - 1) row.status = "pending" row.next_attempt_at = datetime.now(UTC) + timedelta( seconds=RETRY_DELAYS_SECONDS[delay_index] ) else: row.status = "failed" row.next_attempt_at = None await session.commit() async def run_webhook_worker() -> None: logger.info("通话分析 Webhook worker 已启动") while True: try: delivery = await claim_pending_webhook() if delivery: try: result = await deliver_webhook( url=delivery.url, secret=delivery.secret, event_id=delivery.event_id, payload=delivery.payload, ) except UnsafeWebhookUrl as exc: # DNS outages are transient; policy violations are permanent. result = DeliveryResult( False, None if exc.retryable else 400, 0, str(exc) ) await finish_delivery(delivery.id, result) continue except asyncio.CancelledError: raise except Exception as exc: logger.error(f"Webhook worker 暂时不可用:{exc}") await asyncio.sleep(POLL_INTERVAL_SECONDS)