"""FlowDeck — Public API v2 : webhooks. Découpe A28 de l'ancien app/routers/api_v2.py (2 110 lignes, 115 routes) — un module par concern, contrat inchangé (Bearer+scopes, pagination, RFC7807, audit + idempotency). """ from __future__ import annotations import logging from fastapi import APIRouter, Body, Header, HTTPException, Request from fastapi.responses import JSONResponse from app.db import get_conn from app.services.api_v2_helpers import ( # noqa: F401 — require_scope est utilisé par les handlers audit_log, check_idempotency, check_v2_rate_limit, get_bearer_user, has_scope, paginate_headers, parse_pagination, require_scope, row_to_dict, store_idempotency, to_iso8601, validate_scopes_input, ) from ._common import _v2_rate_check logger = logging.getLogger(__name__) router = APIRouter(tags=["api-v2"]) @router.get("/webhooks") def list_webhooks_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit, offset = parse_pagination(request) with get_conn() as conn: total = conn.execute("SELECT COUNT(*) AS n FROM webhook_subscriptions").fetchone()["n"] rows = conn.execute( "SELECT * FROM webhook_subscriptions ORDER BY created_at DESC LIMIT ? OFFSET ?", (limit, offset), ).fetchall() return JSONResponse( content={"webhooks": [dict(r) for r in rows]}, headers=paginate_headers(total), ) @router.post("/webhooks") def create_webhook_v2(request: Request, authorization: str | None = Header(default=None), body: dict = Body(default={})): user = require_scope("write")(request, authorization) _v2_rate_check(request, user) from app.services.webhook_outbound import EVENTS, _event_matches url = (body.get("url") or "").strip() event = (body.get("event") or "page.created").strip() secret = (body.get("secret") or "").strip() if not url or not url.startswith("http"): raise HTTPException(400, "url must start with http") if not event or (event not in EVENTS and not (event.endswith(".*") or event in ("*", "all"))): raise HTTPException(400, f"Unknown event '{event}'. See GET /api/v2/webhooks/events") # make sure the pattern matches at least one known event if not any(_event_matches(event, e) for e in EVENTS): raise HTTPException(400, f"Event pattern '{event}' matches no known event") with get_conn() as conn: cur = conn.execute("INSERT INTO webhook_subscriptions (url, event, secret) VALUES (?, ?, ?)", (url, event, secret)) wid = cur.lastrowid conn.commit() row = conn.execute("SELECT * FROM webhook_subscriptions WHERE id=?", (wid,)).fetchone() audit_log(user, "webhook.create", "webhook", wid, url, request) return {"id": wid, "status": "created", "webhook": dict(row) if row else {}, "signature_header": "X-FlowDeck-Signature (HMAC-SHA256, sha256=)" if secret else None} @router.patch("/webhooks/{webhook_id}") def patch_webhook_v2(webhook_id: int, request: Request, authorization: str | None = Header(default=None), body: dict = Body(default={})): user = require_scope("write")(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: row = conn.execute("SELECT * FROM webhook_subscriptions WHERE id=?", (webhook_id,)).fetchone() if not row: raise HTTPException(404, "Webhook not found") url = body.get("url", row["url"]) event = body.get("event", row["event"]) secret = body.get("secret", row["secret"]) active = int(bool(body.get("active", row["active"]))) conn.execute("UPDATE webhook_subscriptions SET url=?, event=?, secret=?, active=? WHERE id=?", (url, event, secret, active, webhook_id)) conn.commit() audit_log(user, "webhook.update", "webhook", webhook_id, "", request) return {"id": webhook_id, "status": "updated"} @router.delete("/webhooks/{webhook_id}") def delete_webhook_v2(webhook_id: int, request: Request, authorization: str | None = Header(default=None)): user = require_scope("write")(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: conn.execute("DELETE FROM webhook_subscriptions WHERE id=?", (webhook_id,)) conn.commit() audit_log(user, "webhook.delete", "webhook", webhook_id, "", request) return {"id": webhook_id, "status": "deleted"} @router.post("/webhooks/{webhook_id}/test") async def test_webhook_v2(webhook_id: int, request: Request, authorization: str | None = Header(default=None)): user = require_scope("write")(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: row = conn.execute("SELECT * FROM webhook_subscriptions WHERE id=?", (webhook_id,)).fetchone() if not row: raise HTTPException(404, "Webhook not found") # live delivery via the prod dispatcher (HMAC + retry + journal), # direct to this subscription only (no wildcard fan-out) from app.services.webhook_outbound import deliver_to_sub ok = await deliver_to_sub(webhook_id, row["url"], "ping", {"webhook_id": webhook_id, "test": True}, row["secret"] or "") audit_log(user, "webhook.test", "webhook", webhook_id, f"ok={ok}", request) return {"webhook_id": webhook_id, "status": "tested", "delivered": ok} @router.get("/webhooks/{webhook_id}/deliveries") def list_deliveries_v2(webhook_id: int, request: Request, status: str | None = None, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit, offset = parse_pagination(request) with get_conn() as conn: if status: total = conn.execute( "SELECT COUNT(*) AS n FROM webhook_deliveries WHERE webhook_id=? AND status=?", (webhook_id, status)).fetchone()["n"] rows = conn.execute( "SELECT * FROM webhook_deliveries WHERE webhook_id=? AND status=? " "ORDER BY created_at DESC LIMIT ? OFFSET ?", (webhook_id, status, limit, offset)).fetchall() else: total = conn.execute( "SELECT COUNT(*) AS n FROM webhook_deliveries WHERE webhook_id=?", (webhook_id,)).fetchone()["n"] rows = conn.execute( "SELECT * FROM webhook_deliveries WHERE webhook_id=? " "ORDER BY created_at DESC LIMIT ? OFFSET ?", (webhook_id, limit, offset)).fetchall() return JSONResponse( content={"deliveries": [row_to_dict(r) for r in rows]}, headers=paginate_headers(total), ) @router.post("/webhooks/{webhook_id}/retry") async def retry_webhook_deliveries(webhook_id: int, request: Request, authorization: str | None = Header(default=None)): """Manually retry failed deliveries for a webhook.""" user = require_scope("write")(request, authorization) _v2_rate_check(request, user) from app.services.webhook_outbound import retry_due_deliveries with get_conn() as conn: # Force retry by setting next_retry_at to the past conn.execute( """UPDATE webhook_deliveries SET next_retry_at = strftime('%s', 'now', '-1 second') WHERE webhook_id = ? AND status = 'retrying'""", (webhook_id,), ) conn.commit() retried = await retry_due_deliveries() audit_log(user, "webhook.retry", "webhook", webhook_id, f"retried={retried}", request) return {"webhook_id": webhook_id, "status": "retried", "retried_count": retried} @router.post("/webhooks/verify-signature") def verify_webhook_signature(request: Request, authorization: str | None = Header(default=None), body: dict = Body(default={})): """Verify a webhook signature (for debugging/testing).""" user = get_bearer_user(request, authorization) _v2_rate_check(request, user) from app.services.webhook_outbound import verify_signature secret = body.get("secret", "") payload = body.get("payload", "{}") signature = body.get("signature", "") is_valid = verify_signature(secret, payload.encode(), signature) audit_log(user, "webhook.signature_verify", "webhook", 0, f"valid={is_valid}", request) return {"valid": is_valid, "secret": secret[:10] + "..." if len(secret) > 10 else secret}