Files
flowdeck/app/routers/api_v2/webhooks.py
T
bruno 6a5fe0524a
FlowDeck CI / lint (push) Canceled after 0s
FlowDeck CI / test (push) Canceled after 0s
FlowDeck CI / docker (push) Canceled after 0s
refactor: A28 lot 1 — api_v2.py (2 110 L) → package 14 fichiers (v7.29.0)
Découpe par concern de l'ancien app/routers/api_v2.py (2 110 lignes,
115 routes) en package `app/routers/api_v2/` :

- 12 modules de routes : collections 566 L (23 r.), engagement 338 (21),
  workspaces 230 (9), templates_io 205 (9), webhooks 195 (8),
  identity 195 (7), views 164 (8), sharing 160 (8), properties 151 (7),
  planning 148 (7), projects 93 (4), admin 91 (4)
- `_common.py` : helpers partagés (_hash, _v2_rate_check)
- `__init__.py` : router = APIRouter(prefix="/api/v2") + include_router
  sur les routers de sections (sans prefix, tags « api-v2 »)

Preuve contractuelle : `docs/openapi-v2.json` régénéré = IDENTIQUE
byte-à-byte (0 changement de chemin/tag/operation_id). Seul importateur
(app/main.py : from app.routers.api_v2 import router) fonctionne via le
package. En-tête d'imports copié par module puis émondé par ruff --fix
(143 imports morts), I001 réordonnés.

Reste A28 : dashboard.py 2 735 L, collections.py 2 622 L, board.py 2 101 L
(même recette, lots suivants).

suite **1091/1091** · ruff OK · OpenAPI 509 identique · docs à jour
2026-10-01 23:26:45 -04:00

196 lines
8.6 KiB
Python

"""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=<hex>)" 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}