- `app/services/http_client.py` : `async with shared_client(timeout=15) as client:` remplace les 49 créations `async with httpx.AsyncClient(` de 14 fichiers (gitea ×21, providers oidc/oauth ×11, calendar ×4, automations ×3…) — le pool de connexions est réutilisé au lieu d'être recréé à chaque appel. __aexit__ no-op (le client partagé ne se ferme pas à la sortie). - Cache par (boucle d'event, kwargs) en WeakKeyDictionary : un AsyncClient n'est JAMAIS partagé entre deux loops (piège des tests « Event loop is closed ») — une boucle par test = client propre collecté avec la boucle. Clé = kwargs triés, repr() pour les valeurs non hashables (`headers=` dict → TypeError rattrapé par la suite). - Laissés délibérément : github_adapter (transport MockTransport injecté), webhook_outbound (client « own_client » fermé par la fonction). - Tests : `test_http_client_shared_and_loop_scoped` (réutilisation mêmes kwargs / cloisonné kwargs / cloisonné loop) ; le stub des webhooks patche aussi la fabrique `http_client.httpx` + purge du cache (avant : webhook_outbound.httpx patché mais la fabrique partagée créait un vrai client → réseau réel dans les tests). suite **1091/1091** · ruff OK · docs à jour
286 lines
12 KiB
Python
286 lines
12 KiB
Python
"""FlowDeck — Webhook outbound dispatcher (v2.1.0, prod v6.4.0).
|
|
|
|
Production-grade delivery: HMAC-SHA256 signature (``X-FlowDeck-Signature``),
|
|
retry with backoff 2s/10s/60s, full delivery journal in ``webhook_deliveries``.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import hashlib
|
|
import hmac
|
|
import json
|
|
import logging
|
|
import time
|
|
|
|
import httpx
|
|
|
|
from app.db import get_conn
|
|
from app.services.http_client import shared_client
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# ── Event catalogue (≈50 events) ────────────────────────────────────────────
|
|
EVENTS = [
|
|
# pages (block-editor documents)
|
|
"page.created", "page.updated", "page.deleted", "page.moved",
|
|
"page.restored", "page.locked", "page.unlocked",
|
|
"page.shared", "page.published", "page.unpublished",
|
|
"page.renamed", "page.duplicated", "page.archived",
|
|
# comments & collaboration
|
|
"comment.added", "comment.updated", "comment.resolved", "comment.deleted",
|
|
"mention.added", "mention.resolved",
|
|
# collections (databases)
|
|
"collection.created", "collection.updated", "collection.deleted",
|
|
"collection.renamed", "collection.duplicated",
|
|
"collection.page.created", "collection.page.updated", "collection.page.deleted",
|
|
"collection.page.moved", "collection.view.created", "collection.view.updated",
|
|
"collection.view.deleted", "collection.property.created", "collection.property.updated",
|
|
"collection.property.deleted",
|
|
# sprints / tasks
|
|
"sprint.created", "sprint.updated", "sprint.deleted",
|
|
"sprint.started", "sprint.completed", "sprint.canceled",
|
|
# sharing / favorites
|
|
"favorite.added", "favorite.removed",
|
|
"share.created", "share.updated", "share.revoked",
|
|
# automations & agent
|
|
"automation.fired", "automation.failed", "automation.retrying",
|
|
"agent.run.started", "agent.run.finished", "agent.run.failed",
|
|
# workspace & users
|
|
"workspace.created", "workspace.updated", "workspace.deleted",
|
|
"workspace.member.added", "workspace.member.removed", "workspace.member.role_changed",
|
|
# files & imports
|
|
"file.uploaded", "file.deleted",
|
|
"import.started", "import.completed", "import.failed",
|
|
# generic
|
|
"ping",
|
|
]
|
|
|
|
# Retry schedule: delays in seconds between attempts (attempt 0 = immediate).
|
|
RETRY_DELAYS = (2, 10, 60)
|
|
MAX_ATTEMPTS = len(RETRY_DELAYS) + 1 # 1 initial + 3 retries
|
|
|
|
|
|
def sign_payload(secret: str, body: bytes) -> str:
|
|
"""HMAC-SHA256 hex signature of the raw JSON body (``sha256=<hex>``)."""
|
|
return "sha256=" + hmac.new(secret.encode(), body, hashlib.sha256).hexdigest()
|
|
|
|
|
|
def verify_signature(secret: str, body: bytes, signature: str) -> bool:
|
|
"""Constant-time check of an ``X-FlowDeck-Signature`` header value."""
|
|
if not secret or not signature:
|
|
return False
|
|
expected = sign_payload(secret, body)
|
|
return hmac.compare_digest(expected, signature)
|
|
|
|
|
|
def _event_matches(sub_event: str, event: str) -> bool:
|
|
"""Wildcard matching: ``*`` matches all, ``page.*`` matches page.* events."""
|
|
if sub_event in ("*", "all"):
|
|
return True
|
|
if sub_event.endswith(".*"):
|
|
return event.startswith(sub_event[:-1])
|
|
return sub_event == event
|
|
|
|
|
|
def _matching_subs(conn, event: str) -> list:
|
|
rows = conn.execute(
|
|
"SELECT id, url, event, secret FROM webhook_subscriptions WHERE active=1"
|
|
).fetchall()
|
|
return [r for r in rows if _event_matches((r["event"] or "").strip(), event)]
|
|
|
|
|
|
def _log_delivery(webhook_id: int, event: str, status: str, http_code: int | None,
|
|
error: str, duration_ms: int, attempt: int, payload: dict,
|
|
next_retry_at: float | None = None) -> None:
|
|
try:
|
|
with get_conn() as conn:
|
|
conn.execute(
|
|
"""INSERT INTO webhook_deliveries
|
|
(webhook_id, event, status, http_code, error, duration_ms,
|
|
attempt, payload, next_retry_at)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)""",
|
|
(webhook_id, event, status,
|
|
http_code if http_code is not None else 0,
|
|
(error or "")[:1000], duration_ms, attempt,
|
|
json.dumps(payload)[:8000],
|
|
next_retry_at),
|
|
)
|
|
conn.commit()
|
|
except Exception: # noqa: BLE001 — logging must never break delivery
|
|
logger.debug("Failed to log webhook delivery", exc_info=True)
|
|
|
|
|
|
async def _deliver_once(client: httpx.AsyncClient, url: str, event: str,
|
|
payload: dict, secret: str) -> tuple[int | None, str]:
|
|
"""Single POST attempt. Returns (http_code, error).
|
|
|
|
Includes HMAC-SHA256 signature in X-FlowDeck-Signature header.
|
|
"""
|
|
body = json.dumps({"event": event, **payload}).encode()
|
|
headers = {
|
|
"Content-Type": "application/json",
|
|
"X-FlowDeck-Event": event,
|
|
"User-Agent": "FlowDeck-Webhooks/1.0",
|
|
}
|
|
if secret:
|
|
headers["X-FlowDeck-Signature"] = sign_payload(secret, body)
|
|
# legacy header kept for backward compatibility
|
|
headers["X-FlowDeck-Secret"] = secret
|
|
try:
|
|
resp = await client.post(url, content=body, headers=headers, timeout=30.0)
|
|
if 200 <= resp.status_code < 300:
|
|
return resp.status_code, ""
|
|
return resp.status_code, f"HTTP {resp.status_code}"
|
|
except httpx.TimeoutException as exc:
|
|
return None, f"Timeout: {str(exc)[:400]}"
|
|
except httpx.RequestError as exc:
|
|
return None, f"Request error: {str(exc)[:400]}"
|
|
except Exception as exc: # noqa: BLE001
|
|
return None, str(exc)[:500]
|
|
|
|
|
|
async def deliver_to_sub(sub_id: int, url: str, event: str, payload: dict,
|
|
secret: str, *, _client: httpx.AsyncClient | None = None) -> bool:
|
|
"""Deliver with retry (immediate + 2s/10s/60s). Logs every attempt.
|
|
|
|
Returns True on success. Failures are re-queued via ``next_retry_at`` so
|
|
the background scheduler can pick them up even if this process restarts.
|
|
|
|
Retry schedule: 2s, 10s, 60s (3 retries total + initial attempt).
|
|
"""
|
|
own_client = _client is None
|
|
client = _client or httpx.AsyncClient(timeout=30.0)
|
|
try:
|
|
for attempt in range(MAX_ATTEMPTS):
|
|
if attempt > 0:
|
|
delay = RETRY_DELAYS[attempt - 1]
|
|
logger.debug("Webhook retry attempt %d/%d for %s after %ds", attempt, MAX_ATTEMPTS, url, delay)
|
|
await asyncio.sleep(delay)
|
|
start = time.monotonic()
|
|
code, err = await _deliver_once(client, url, event, payload, secret or "")
|
|
duration = int((time.monotonic() - start) * 1000)
|
|
ok = code is not None and 200 <= code < 300
|
|
last = attempt == MAX_ATTEMPTS - 1
|
|
if ok or last:
|
|
_log_delivery(sub_id, event, "delivered" if ok else "failed",
|
|
code, err, duration, attempt, payload)
|
|
return ok
|
|
# schedule retry
|
|
delay = RETRY_DELAYS[attempt]
|
|
_log_delivery(sub_id, event, "retrying", code, err, duration,
|
|
attempt, payload,
|
|
next_retry_at=time.time() + delay)
|
|
return False
|
|
finally:
|
|
if own_client:
|
|
await client.aclose()
|
|
|
|
|
|
async def fire_event(event: str, payload: dict):
|
|
"""Fire a webhook event to all registered subscribers (wildcard-aware).
|
|
|
|
Unknown events are ignored. Delivery itself never raises: each
|
|
subscription is attempted independently and journaled.
|
|
"""
|
|
if event not in EVENTS:
|
|
logger.debug("Ignoring unknown webhook event: %s", event)
|
|
return
|
|
with get_conn() as conn:
|
|
subs = _matching_subs(conn, event)
|
|
if not subs:
|
|
logger.debug("No subscribers for event: %s", event)
|
|
return
|
|
async with shared_client(timeout=30.0) as client:
|
|
for sub in subs:
|
|
try:
|
|
await deliver_to_sub(sub["id"], sub["url"], event,
|
|
dict(payload), sub["secret"] or "",
|
|
_client=client)
|
|
except Exception as exc: # noqa: BLE001
|
|
logger.error("Webhook delivery failed to %s: %s", sub["url"], exc)
|
|
|
|
|
|
async def retry_due_deliveries(now: float | None = None) -> int:
|
|
"""Re-fire deliveries stuck in ``retrying`` whose ``next_retry_at`` passed.
|
|
|
|
Returns the number of deliveries retried. Called by the background
|
|
scheduler and directly testable.
|
|
"""
|
|
now = now if now is not None else time.time()
|
|
with get_conn() as conn:
|
|
rows = conn.execute(
|
|
"""SELECT d.id, d.webhook_id, d.event, d.payload, d.attempt,
|
|
s.url, s.secret
|
|
FROM webhook_deliveries d
|
|
JOIN webhook_subscriptions s ON s.id = d.webhook_id
|
|
WHERE d.status='retrying' AND s.active=1
|
|
AND d.next_retry_at IS NOT NULL AND d.next_retry_at <= ?
|
|
ORDER BY d.next_retry_at LIMIT 50""",
|
|
(now,),
|
|
).fetchall()
|
|
if not rows:
|
|
return 0
|
|
async with shared_client(timeout=10) as client:
|
|
for r in rows:
|
|
try:
|
|
payload = json.loads(r["payload"] or "{}")
|
|
if not isinstance(payload, dict):
|
|
payload = {"data": payload}
|
|
except Exception:
|
|
payload = {}
|
|
event = r["event"] or "ping"
|
|
# resume where the previous run stopped, up to MAX_ATTEMPTS
|
|
start_attempt = min((r["attempt"] or 0) + 1, MAX_ATTEMPTS - 1)
|
|
for att in range(start_attempt, MAX_ATTEMPTS):
|
|
delay = RETRY_DELAYS[att - 1] if att > 0 else 0
|
|
if delay:
|
|
await asyncio.sleep(delay)
|
|
t0 = time.monotonic()
|
|
code, err = await _deliver_once(client, r["url"], event,
|
|
payload, r["secret"] or "")
|
|
duration = int((time.monotonic() - t0) * 1000)
|
|
ok = code is not None and 200 <= code < 300
|
|
last = att == MAX_ATTEMPTS - 1
|
|
if ok or last:
|
|
_log_delivery(r["webhook_id"], event,
|
|
"delivered" if ok else "failed",
|
|
code, err, duration, att, payload)
|
|
break
|
|
_log_delivery(r["webhook_id"], event, "retrying", code, err,
|
|
duration, att, payload,
|
|
next_retry_at=time.time() + RETRY_DELAYS[att])
|
|
with get_conn() as conn:
|
|
conn.execute("UPDATE webhook_deliveries SET status='superseded' WHERE id=?", (r["id"],))
|
|
conn.commit()
|
|
return len(rows)
|
|
|
|
|
|
async def webhook_retry_scheduler(interval_seconds: int = 60) -> None:
|
|
"""Background loop: retry due webhook deliveries every minute."""
|
|
from app.config import settings as _settings
|
|
interval = getattr(_settings, "webhook_retry_interval_seconds", interval_seconds)
|
|
while True:
|
|
try:
|
|
await asyncio.sleep(interval)
|
|
await retry_due_deliveries()
|
|
except asyncio.CancelledError:
|
|
raise
|
|
except Exception: # noqa: BLE001
|
|
logger.debug("webhook_retry_scheduler tick failed", exc_info=True)
|
|
|
|
|
|
def init_webhook_tables():
|
|
"""Create the webhook_subscriptions table if it doesn't exist."""
|
|
with get_conn() as conn:
|
|
conn.execute("""
|
|
CREATE TABLE IF NOT EXISTS webhook_subscriptions (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
url TEXT NOT NULL,
|
|
event TEXT NOT NULL,
|
|
secret TEXT DEFAULT '',
|
|
active BOOLEAN NOT NULL DEFAULT 1,
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
|
|
)
|
|
""")
|
|
conn.commit()
|