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
196 lines
8.6 KiB
Python
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}
|