Spaces:
Runtime error
Runtime error
Download api_webhooks.py from sdudeja/agentic-extractor: direct link, hf CLI and curl.
- Browser
- Download file 6.9 kB
-
https://huggingface.co/spaces/sdudeja/agentic-extractor/resolve/main/api_webhooks.py
- Command line
-
hf download hf://spaces/sdudeja/agentic-extractor/api_webhooks.py
-
curl -L -o api_webhooks.py https://huggingface.co/spaces/sdudeja/agentic-extractor/resolve/main/api_webhooks.py
6.9 kB
| """ | |
| DocuLens — Webhook management and delivery. | |
| Allows users to register webhook URLs that receive POST notifications | |
| when extraction or validation completes. | |
| """ | |
| import os | |
| import json | |
| import hmac | |
| import hashlib | |
| import logging | |
| import time | |
| from typing import Optional | |
| from uuid import uuid4 | |
| import requests | |
| from fastapi import APIRouter, HTTPException, Depends | |
| from pydantic import BaseModel, HttpUrl | |
| from middleware import check_rate_limit | |
| logger = logging.getLogger(__name__) | |
| webhooks_router = APIRouter() | |
| # --------------------------------------------------------------------------- | |
| # In-memory webhook store (Supabase-backed when available) | |
| # --------------------------------------------------------------------------- | |
| _webhooks: dict[str, dict] = {} # keyed by webhook_id | |
| def _get_webhooks_from_db() -> list[dict]: | |
| """Load webhooks from Supabase if available.""" | |
| try: | |
| from db.supabase import get_client | |
| client = get_client() | |
| if not client: | |
| return [] | |
| resp = client.table("webhooks").select("*").eq("is_active", True).execute() | |
| return resp.data or [] | |
| except Exception as e: | |
| logger.debug(f"Could not load webhooks from DB: {e}") | |
| return [] | |
| def _save_webhook_to_db(webhook: dict) -> bool: | |
| """Persist webhook to Supabase.""" | |
| try: | |
| from db.supabase import get_client | |
| client = get_client() | |
| if not client: | |
| return False | |
| client.table("webhooks").insert(webhook).execute() | |
| return True | |
| except Exception as e: | |
| logger.warning(f"Failed to save webhook to DB: {e}") | |
| return False | |
| def _delete_webhook_from_db(webhook_id: str) -> bool: | |
| """Remove webhook from Supabase.""" | |
| try: | |
| from db.supabase import get_client | |
| client = get_client() | |
| if not client: | |
| return False | |
| client.table("webhooks").delete().eq("id", webhook_id).execute() | |
| return True | |
| except Exception as e: | |
| logger.warning(f"Failed to delete webhook from DB: {e}") | |
| return False | |
| # --------------------------------------------------------------------------- | |
| # Models | |
| # --------------------------------------------------------------------------- | |
| class WebhookCreate(BaseModel): | |
| url: HttpUrl | |
| events: list[str] = ["extraction.completed", "validation.completed"] | |
| secret: Optional[str] = None # optional signing secret | |
| class WebhookResponse(BaseModel): | |
| id: str | |
| url: str | |
| events: list[str] | |
| is_active: bool | |
| created_at: str | |
| # --------------------------------------------------------------------------- | |
| # Endpoints | |
| # --------------------------------------------------------------------------- | |
| async def register_webhook( | |
| body: WebhookCreate, | |
| api_key: str = Depends(check_rate_limit), | |
| ): | |
| """Register a webhook URL to receive extraction/validation result notifications.""" | |
| valid_events = {"extraction.completed", "validation.completed", "batch.completed"} | |
| for event in body.events: | |
| if event not in valid_events: | |
| raise HTTPException( | |
| status_code=400, | |
| detail=f"Invalid event '{event}'. Valid events: {sorted(valid_events)}", | |
| ) | |
| webhook_id = str(uuid4()) | |
| webhook = { | |
| "id": webhook_id, | |
| "url": str(body.url), | |
| "events": body.events, | |
| "secret": body.secret or "", | |
| "is_active": True, | |
| "created_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()), | |
| } | |
| _save_webhook_to_db(webhook) | |
| _webhooks[webhook_id] = webhook | |
| return { | |
| "id": webhook_id, | |
| "url": webhook["url"], | |
| "events": webhook["events"], | |
| "is_active": True, | |
| "created_at": webhook["created_at"], | |
| } | |
| async def list_webhooks(api_key: str = Depends(check_rate_limit)): | |
| """List all registered webhooks.""" | |
| # Merge in-memory and DB webhooks | |
| db_hooks = _get_webhooks_from_db() | |
| all_hooks = {h["id"]: h for h in db_hooks} | |
| all_hooks.update(_webhooks) | |
| return { | |
| "webhooks": [ | |
| { | |
| "id": h["id"], | |
| "url": h["url"], | |
| "events": h.get("events", []), | |
| "is_active": h.get("is_active", True), | |
| "created_at": h.get("created_at", ""), | |
| } | |
| for h in all_hooks.values() | |
| if h.get("is_active", True) | |
| ] | |
| } | |
| async def delete_webhook(webhook_id: str, api_key: str = Depends(check_rate_limit)): | |
| """Deactivate a webhook.""" | |
| if webhook_id in _webhooks: | |
| _webhooks[webhook_id]["is_active"] = False | |
| _delete_webhook_from_db(webhook_id) | |
| return {"deleted": True} | |
| # --------------------------------------------------------------------------- | |
| # Delivery — called internally after extraction/validation completes | |
| # --------------------------------------------------------------------------- | |
| WEBHOOK_TIMEOUT = int(os.environ.get("WEBHOOK_TIMEOUT", "10")) # seconds | |
| def _sign_payload(payload_bytes: bytes, secret: str) -> str: | |
| """Create HMAC-SHA256 signature for webhook payload.""" | |
| return hmac.new(secret.encode(), payload_bytes, hashlib.sha256).hexdigest() | |
| def deliver_webhook(event: str, data: dict): | |
| """ | |
| Fire webhooks for a given event. Called internally after | |
| extraction or validation completes. Non-blocking best-effort. | |
| """ | |
| # Gather all active webhooks (in-memory + DB) | |
| all_hooks = list(_webhooks.values()) | |
| try: | |
| all_hooks.extend(_get_webhooks_from_db()) | |
| except Exception: | |
| pass | |
| # Deduplicate by ID | |
| seen = set() | |
| unique_hooks = [] | |
| for h in all_hooks: | |
| if h["id"] not in seen and h.get("is_active", True): | |
| seen.add(h["id"]) | |
| unique_hooks.append(h) | |
| for hook in unique_hooks: | |
| if event not in hook.get("events", []): | |
| continue | |
| payload = { | |
| "event": event, | |
| "timestamp": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()), | |
| "data": data, | |
| } | |
| payload_bytes = json.dumps(payload, default=str).encode() | |
| headers = {"Content-Type": "application/json"} | |
| if hook.get("secret"): | |
| headers["X-Webhook-Signature"] = _sign_payload(payload_bytes, hook["secret"]) | |
| try: | |
| resp = requests.post( | |
| hook["url"], | |
| data=payload_bytes, | |
| headers=headers, | |
| timeout=WEBHOOK_TIMEOUT, | |
| ) | |
| logger.info( | |
| "Webhook delivered: event=%s url=%s status=%d", | |
| event, hook["url"], resp.status_code, | |
| ) | |
| except Exception as e: | |
| logger.warning("Webhook delivery failed: url=%s error=%s", hook["url"], e) | |