Windpaint
Recipes

Webhook receiver

A FastAPI service that receives per-job webhooks, confirms each one against the API, saves the outputs exactly once, and catches jobs whose webhook never arrived.

When you submit a job with webhook_url, Windpaint POSTs the job’s status to that URL once the job finishes. This recipe builds a receiver for those calls: it answers immediately, re-fetches the job’s status with your key, downloads the outputs once per request_id, and runs a small reconcile loop for jobs whose webhook was lost.

Why the receiver doesn’t trust the payload: per-job webhooks are not signed, are sent once with a 10 second timeout, and are not retried. Anyone who learns your URL can POST to it, and a delivery that fails is gone. So the webhook is a hint that a job may be done. The GET /v1/generation/requests/{id}/status call, made with your key, is the source of truth, and a spoofed call can’t fake its answer.

Prerequisites

  • An API key exported as WINDPAINT_API_KEY (API keys).
  • Python 3.10+ with pip install fastapi uvicorn httpx.
  • A tunnel for local testing: cloudflared or ngrok. Webhooks are only sent to https:// URLs.

The code

Three files: a small SQLite store shared by both sides, the receiver, and a submit script.

store.py: track every job by request_id

store.py
import sqlite3

db = sqlite3.connect("jobs.db", isolation_level=None, check_same_thread=False)
db.execute(
    """CREATE TABLE IF NOT EXISTS jobs (
        request_id TEXT PRIMARY KEY,
        state      TEXT NOT NULL,      -- pending | processing | completed | failed | nsfw | canceled | error
        updated_at TEXT NOT NULL DEFAULT (datetime('now'))
    )"""
)


def record_pending(request_id: str) -> None:
    db.execute("INSERT OR IGNORE INTO jobs (request_id, state) VALUES (?, 'pending')", (request_id,))


def claim(request_id: str) -> bool:
    """Atomically take a job for processing. False if it's done or someone else has it."""
    record_pending(request_id)
    cursor = db.execute(
        "UPDATE jobs SET state = 'processing', updated_at = datetime('now') "
        "WHERE request_id = ? AND state IN ('pending', 'error')",
        (request_id,),
    )
    return cursor.rowcount == 1


def finish(request_id: str, state: str) -> None:
    db.execute(
        "UPDATE jobs SET state = ?, updated_at = datetime('now') WHERE request_id = ?", (state, request_id)
    )


def pending_older_than(seconds: int) -> list[str]:
    rows = db.execute(
        "SELECT request_id FROM jobs WHERE state IN ('pending', 'error') "
        "AND updated_at < datetime('now', ?)",
        (f"-{seconds} seconds",),
    )
    return [row[0] for row in rows]

claim() is what makes the receiver idempotent: only one caller can move a job from pending to processing, so a duplicate webhook, a spoofed one, and the reconcile loop can’t download the same outputs twice.

app.py: the receiver

app.py
import asyncio
import mimetypes
import os
import re
from contextlib import asynccontextmanager
from pathlib import Path

import httpx
from fastapi import BackgroundTasks, FastAPI, Request, Response

import store

API = os.environ.get("WINDPAINT_API_URL", "https://api.windpaint.ai").rstrip("/") + "/v1"
KEY = os.environ["WINDPAINT_API_KEY"]
OUT = Path("outputs")
OUT.mkdir(exist_ok=True)
TERMINAL = {"completed", "failed", "nsfw", "canceled"}
UUID = re.compile(r"^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$")


async def process(request_id: str) -> None:
    if not store.claim(request_id):
        return  # already handled, or being handled
    try:
        async with httpx.AsyncClient(
            base_url=API, headers={"Authorization": f"Bearer {KEY}"}, timeout=60
        ) as api:
            response = await api.get(f"/generation/requests/{request_id}/status")
            if response.status_code == 404:
                store.finish(request_id, "ignored")  # not a job in your organization
                return
            response.raise_for_status()
            status = response.json()

            if status["status"] not in TERMINAL:
                store.finish(request_id, "pending")  # early or spoofed call; reconcile will retry
                return

            if status["status"] == "completed":
                for i, output in enumerate(status["outputs"]):
                    signed = await api.get(f"/generation/assets/{output['id']}/url")
                    signed.raise_for_status()
                    url = signed.json()["data"]["url"]
                    async with httpx.AsyncClient(timeout=300) as plain:  # no Authorization on signed URLs
                        media = await plain.get(url)
                        media.raise_for_status()
                    ext = mimetypes.guess_extension(output["content_type"]) or ""
                    (OUT / f"{request_id}-{i}{ext}").write_bytes(media.content)
                print(f"{request_id}: saved {len(status['outputs'])} output(s), "
                      f"{status['credits']['actual']} credits")
            else:
                print(f"{request_id}: {status['status']} ({status['error']})")
            store.finish(request_id, status["status"])
    except Exception:
        store.finish(request_id, "error")  # picked up again by the reconcile loop
        raise


async def reconcile_forever() -> None:
    """Catch jobs whose webhook never arrived: webhooks are sent once and not retried."""
    while True:
        await asyncio.sleep(60)
        for request_id in store.pending_older_than(120):
            try:
                await process(request_id)
            except Exception as exc:
                print(f"{request_id}: reconcile failed: {exc!r}")


@asynccontextmanager
async def lifespan(app: FastAPI):
    task = asyncio.create_task(reconcile_forever())
    yield
    task.cancel()


app = FastAPI(lifespan=lifespan)


@app.post("/windpaint/webhook")
async def webhook(request: Request, background: BackgroundTasks) -> Response:
    try:
        payload = await request.json()
    except ValueError:
        return Response(status_code=400)
    request_id = payload.get("request_id") if isinstance(payload, dict) else None
    if not isinstance(request_id, str) or not UUID.match(request_id):
        return Response(status_code=400)
    # Answer within the 10 s timeout; do the work after responding.
    background.add_task(process, request_id)
    return Response(status_code=204)

The handler reads one field from the body, request_id, and ignores the rest. Everything it acts on comes from the status call. The status endpoint finds a job anywhere in your organization, so a 404 means the id isn’t yours.

submit.py: submit a job with webhook_url

submit.py
import os
import sys

import httpx

import store

API = os.environ.get("WINDPAINT_API_URL", "https://api.windpaint.ai").rstrip("/") + "/v1"

response = httpx.post(
    f"{API}/generation/capabilities/image.generate",
    headers={"Authorization": f"Bearer {os.environ['WINDPAINT_API_KEY']}"},
    json={
        "prompt": sys.argv[1],
        "aspect_ratio": "16:9",
        "webhook_url": os.environ["WEBHOOK_URL"],
    },
    timeout=60,
)
response.raise_for_status()
job = response.json()
store.record_pending(job["request_id"])  # the reconcile loop now owns this job too
print(f"submitted {job['request_id']} ({job['credits_estimate']} credits held)")

Recording the job at submit time is what lets the reconcile loop find it if the webhook never comes. In your own app, this is the row you already store next to the user’s request.

Test it locally

Start the receiver

export WINDPAINT_API_KEY=aak_...
uvicorn app:app --port 8000

Expose it over https

In a second terminal, with either tool:

cloudflared
cloudflared tunnel --url http://localhost:8000
# prints https://<random-words>.trycloudflare.com

Submit a job

In a third terminal, from the same directory (so it shares jobs.db):

export WINDPAINT_API_KEY=aak_...
export WEBHOOK_URL=https://<your-tunnel-host>/windpaint/webhook
python submit.py "a paper boat on a rain-soaked street, reflections, night"

A few seconds later the receiver logs:

INFO:     34.x.x.x:0 - "POST /windpaint/webhook HTTP/1.1" 204 No Content
7c1e9a02-...: saved 1 output(s), 0.08 credits

and the image is in outputs/.

Check the failure paths

  • Duplicate: replay the call with curl -X POST "$WEBHOOK_URL" -H 'Content-Type: application/json' -d '{"request_id": "<id>"}'. Nothing is downloaded again.
  • Lost webhook: stop uvicorn, submit a job, start uvicorn again. Within about three minutes the reconcile loop fetches the job and saves its output.
  • Spoof: POST a random UUID. The status call returns 404 and the id is marked ignored.

Tunnel URLs change each time you restart the tunnel. A job carries the webhook_url it was submitted with, so jobs submitted against an old tunnel URL are only picked up by the reconcile loop.

When the webhook fires

Job outcomeWebhook sent?
completed, failed, nsfwYes, once
Canceled with POST /v1/generation/requests/{id}/cancelNo
Canceled while a render was already runningPossibly, when the render returns

Product runs have no webhooks; poll GET /v1/workflows/runs/{id}. The signed org webhooks cover account events (members, keys, low balance), not generation.

In production

  • Serve the receiver over https://. Webhooks to http:// URLs aren’t sent.
  • Use a path that’s hard to guess if you like, but don’t rely on it for security. The re-fetch is what makes the receiver safe.
  • Swap SQLite for your database. Keep the conditional update in claim(): it’s what keeps two workers from processing the same job.
  • Keep the reconcile loop. It’s the only thing that catches a webhook lost to a deploy, a timeout or a network blip.
  • Store asset ids rather than copying files if you don’t need them on disk; outputs stay in your project. See Media in your app.

Next steps