121 lines
4.0 KiB
Python
121 lines
4.0 KiB
Python
"""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)
|