Files
leocrm/app/routes/webhooks.py
T
Agent Zero abbe7a18fc fix(audit): P0-P3 audit fixes — 838 ruff errors → 0, 30 F821 bugs fixed, 118 files changed
- P0: hooks.py 3-tuple fix, trigger_dispatcher Contract, contacts/plugin unregister_actions_by_owner
- P0: 5 test files — check_permission mocks removed, hardcoded DB credential → env var
- P1: attachment_service DmsFile via Contract helper, restore_registry/history_hooks dedup
- P1: mail/plugin restore unregister, mcp_client datetime.now(UTC), saved_views/filters patterns
- P1: ProtectedRoute fail-closed, 13 test assertion fixes (bcrypt, DB-URLs, SECRET_KEYs)
- P2: deprecated notifications → post_system_message (3 files), forgejo Base, report_generator lazy import
- P2: webhooks permissions, deps.py/roles.py plugin perms removed, import_export default
- P2: address/tags/entity_links patterns removed, worker.py Contract-Umgehungen fixed
- P2: 28 frontend TODOs (hardcoded constants, deprecated notification API)
- P3: dead code, duplicates, deprecated imports, private attr, __import__ inline
- P3: 8 frontend TODOs (LucideIcons, inline styles, XSS, i18n)
- ruff: 838 → 0 (612 auto-fix + 246 manual + 27 F821 regression fix)
- F821: 30 → 0 (AutomationDefinition, DmsFile, user_id, Path, Any, String)
- Contract-Umgehungen: 2 neue gefunden (worker.py:169, worker.py:280) und gefixt
2026-08-16 01:17:18 +02:00

207 lines
7.0 KiB
Python

"""API routes for Webhook CRUD and testing."""
from __future__ import annotations
import uuid
from fastapi import APIRouter, Depends, HTTPException, Query, Request
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.db import get_db
from app.core.rate_limit import RateLimitPolicy, check_rate_limit_policy
from app.deps import get_current_user, require_permission
from app.schemas.webhook import (
WebhookCreate,
WebhookResponse,
WebhookUpdate,
)
from app.services import webhook_service
router = APIRouter(prefix="/api/v1/webhooks", tags=["webhooks"])
@router.get(
"",
response_model=list[WebhookResponse],
dependencies=[Depends(require_permission("workflows:read"))],
)
async def list_webhooks(
event: str | None = Query(None, description="Filter by event name"),
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(get_current_user),
):
"""List all webhooks for the current tenant."""
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
is_admin = current_user.get("is_system_admin", False)
try:
webhooks = await webhook_service.list_webhooks(db, tenant_id, event=event, user_id=user_id, is_system_admin=is_admin)
return webhooks
except PermissionError as e:
raise HTTPException(status_code=403, detail=str(e)) from e
@router.post(
"",
response_model=WebhookResponse,
status_code=201,
dependencies=[Depends(require_permission("workflows:write"))],
)
async def create_webhook(
body: WebhookCreate,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(get_current_user),
):
"""Create a new webhook subscription."""
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
is_admin = current_user.get("is_system_admin", False)
try:
webhook = await webhook_service.create_webhook(
db, tenant_id, user_id, body.model_dump(), is_system_admin=is_admin
)
return webhook
except PermissionError as e:
raise HTTPException(status_code=403, detail=str(e)) from e
@router.get(
"/{webhook_id}",
response_model=WebhookResponse,
dependencies=[Depends(require_permission("automation:read"))],
)
async def get_webhook(
webhook_id: str,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(get_current_user),
):
"""Get a single webhook by ID."""
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
is_admin = current_user.get("is_system_admin", False)
try:
wh_id = uuid.UUID(webhook_id)
except (ValueError, TypeError):
raise HTTPException(400, detail={"detail": "Invalid webhook_id", "code": "invalid_id"}) from None
try:
webhook = await webhook_service.get_webhook(db, tenant_id, wh_id, user_id=user_id, is_system_admin=is_admin)
except PermissionError as e:
raise HTTPException(status_code=403, detail=str(e)) from e
if webhook is None:
raise HTTPException(404, detail={"detail": "Webhook not found", "code": "not_found"})
return webhook
@router.patch(
"/{webhook_id}",
response_model=WebhookResponse,
dependencies=[Depends(require_permission("automation:write"))],
)
async def update_webhook(
webhook_id: str,
body: WebhookUpdate,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(get_current_user),
):
"""Update an existing webhook subscription."""
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
is_admin = current_user.get("is_system_admin", False)
try:
wh_id = uuid.UUID(webhook_id)
except (ValueError, TypeError):
raise HTTPException(400, detail={"detail": "Invalid webhook_id", "code": "invalid_id"}) from None
update_data = {k: v for k, v in body.model_dump().items() if v is not None}
if not update_data:
raise HTTPException(400, detail={"detail": "No fields to update", "code": "no_updates"})
try:
webhook = await webhook_service.update_webhook(
db, tenant_id, wh_id, update_data, user_id=user_id, is_system_admin=is_admin
)
except PermissionError as e:
raise HTTPException(status_code=403, detail=str(e)) from e
if webhook is None:
raise HTTPException(404, detail={"detail": "Webhook not found", "code": "not_found"})
return webhook
@router.delete(
"/{webhook_id}",
status_code=204,
dependencies=[Depends(require_permission("automation:write"))],
)
async def delete_webhook(
webhook_id: str,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(get_current_user),
):
"""Delete a webhook subscription."""
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
is_admin = current_user.get("is_system_admin", False)
try:
wh_id = uuid.UUID(webhook_id)
except (ValueError, TypeError):
raise HTTPException(400, detail={"detail": "Invalid webhook_id", "code": "invalid_id"}) from None
try:
deleted = await webhook_service.delete_webhook(db, tenant_id, wh_id, user_id=user_id, is_system_admin=is_admin)
except PermissionError as e:
raise HTTPException(status_code=403, detail=str(e)) from e
if not deleted:
raise HTTPException(404, detail={"detail": "Webhook not found", "code": "not_found"})
return None
@router.post(
"/{webhook_id}/test",
dependencies=[Depends(require_permission("automation:write"))],
)
async def test_webhook(
webhook_id: str,
request: Request,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(get_current_user),
):
"""Send a test payload to a webhook to verify connectivity."""
tenant_id = uuid.UUID(current_user["tenant_id"])
user_id = uuid.UUID(current_user["user_id"])
is_admin = current_user.get("is_system_admin", False)
# Rate limit — WEBHOOK policy
await check_rate_limit_policy(
f"rate:webhook:test:{tenant_id}:{user_id}",
RateLimitPolicy.WEBHOOK,
)
try:
wh_id = uuid.UUID(webhook_id)
except (ValueError, TypeError):
raise HTTPException(400, detail={"detail": "Invalid webhook_id", "code": "invalid_id"}) from None
try:
webhook = await webhook_service.get_webhook(db, tenant_id, wh_id, user_id=user_id, is_system_admin=is_admin)
except PermissionError as e:
raise HTTPException(status_code=403, detail=str(e)) from e
if webhook is None:
raise HTTPException(404, detail={"detail": "Webhook not found", "code": "not_found"})
test_payload = {
"type": "test",
"message": "This is a test webhook from LeoCRM",
"webhook_id": str(webhook.id),
}
try:
result = await webhook_service.send_webhook(webhook, "webhook.test", test_payload, user_id=user_id, is_system_admin=is_admin)
except PermissionError as e:
raise HTTPException(status_code=403, detail=str(e)) from e
return result