Files
ai-video-fullstack/backend/services/webhooks/worker.py
2026-08-07 14:08:49 +08:00

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)