Soumabrata

Engineering 12 min read May 2026

Building a Production-Grade Webhook Relay
Asynchronous Event-Driven Architecture with Python & Kafka

How to build a self-hosted webhook relay that normalizes, transforms, and reliably delivers events with circuit breaking, rate limiting, exponential backoff, and full audit logging — using FastAPI, Kafka, PostgreSQL, and Redis.

1. The Problem

Webhooks are the glue of modern SaaS — Stripe sends payment events, GitHub pushes repository notifications, and Slack fires interactions. But each provider speaks a different dialect: different signature schemes, payload shapes, and delivery guarantees.

Every webhook integration is bespoke. You write a handler for Stripe, another for GitHub, a third for Slack. The handler must validate the signature, deduplicate retries, transform the payload to your internal schema, make an HTTP request to your backend, and retry if that request fails. Do this for N providers and you've written the same boilerplate N times.

A webhook relay service centralizes this:

Provider → Webhook Relay → Transform Pipeline → Delivery to N Destinations

One endpoint per provider, zero code per integration. Configure routes, filters, and transforms at runtime.

2. Architecture Overview

The relay is a modular monolith — three Python processes sharing a codebase:

Process Role Stack
Gateway (FastAPI) Ingestion, validation, persistence asyncpg, redis-py, aiokafka
Transform Worker Route resolution, filtering, transformation aiokafka, jmespath, custom template engine
Delivery Worker HTTP delivery, retries, circuit breaker, rate limiter httpx, redis-py, aiokafka

Data flows through three Kafka topics:

raw-events → transformed-events → dead-letter

Webhook Source (Stripe, GitHub…) Ingestion Gateway FastAPI · /hooks/{id} PostgreSQL Events · Routes · Users Redis Idempotency · CB · Rate Limiter K A F K A raw-events transformed-events dead-letter Transform Worker Filter · JMESPath · Templates Delivery Worker HTTP · Retry · CB · DLQ Destinations (Slack, your API…) POST /hooks/{id} Save event Publish Consume Consume DLQ on exhaustion HTTP delivery Delivery reads CB/rl state Kafka broker (KRaft, no ZooKeeper) 3 topics · Consumer groups for scaling
End-to-end message flow: webhook source → gateway → Kafka → transform → Kafka → delivery → destination

The critical design decision: the gateway never talks to the workers directly. It writes to PostgreSQL (for durability) and publishes a lightweight pointer to Kafka (for decoupling). The 202 Accepted response is sent the moment the event hits the database, not when it reaches the destination.

3. Ingestion Gateway: Where Webhooks Land

The gateway endpoint POST /hooks/{endpoint_id} has a tight contract: validate fast, persist, respond.

# Pseudocode flow
1. Read body (streaming, capped at 1 MB)
2. Look up endpoint from DB
3. Check IP allowlist
4. Verify HMAC-SHA256 signature (if configured)
5. Check Redis idempotency key (24h TTL)
6. Save event + 'pending' to PostgreSQL
7. Fire-and-forget Kafka publish → 'queued' or 'failed'
8. Return 202 Accepted

Idempotency

Uses Redis SET NX EX — if the same Idempotency-Key header arrives twice, the second request gets back {"status": "duplicate"}. This prevents accidental double-charges from Stripe retries.

Signature Verification

Constant-time comparison (hmac.compare_digest) with an explicit sha256= prefix parse. No surprises when a provider changes their scheme.

Event Status Tracking

Events start as pending. A fire-and-forget task publishes to Kafka and updates the status to queued on success or failed on failure. This makes Kafka downtime observable rather than silent data loss — events stuck at pending can be replayed via the API later.

4. Kafka as the Decoupling Layer

Why Kafka and not Redis pub/sub or RabbitMQ?

The trade-off: Kafka is stateful. If it's down, the gateway shouldn't crash. Graceful degradation is explicit:

async def init_kafka() → AIOKafkaProducer | None:
    _producer = AIOKafkaProducer(...)
    try:
        await _producer.start()
    except Exception:
        logger.warning("Kafka unreachable — "
                       "continuing without producer")
        await _producer.stop()
        _producer = None
    return _producer

Every caller checks if producer is None. The transform and delivery workers won't run, but the gateway still accepts webhooks, validates them, and saves them to PostgreSQL. When Kafka recovers, POST /api/events/{id}/replay re-publishes the event.

5. The Transform Pipeline: Payload as Code

Each route can define a transform_pipeline — an ordered list of steps applied to the payload before delivery:

[
  {"type": "jmespath", "expression": "data.object"},
  {"type": "template", "body": {
    "amount": "{{amount / 100}}",
    "currency": "{{currency}}"
  }}
]
Step Type Description
passthrough No-op (useful as placeholder)
jmespath JMESPath expression to select/restructure data
template Substitutes {{expr}} with dot-path access and basic arithmetic

The template engine avoids eval() — only key traversal and division:

# "data.amount / 100" → payload["data"]["amount"] / 100
parts = expr.split("/")
value = traverse(payload, parts[0])
if len(parts) > 1:
    value = value / float(parts[1])  # if numeric

Route Filtering

The filter_expression field on a route is a JMESPath expression evaluated against a unified context:

context = {
    "body":    payload,
    "headers": {"x-github-event": "push"}
}

This enables filters like headers."x-github-event" == 'push' or body.event_type == 'payment.succeeded'. Routes whose filter doesn't match are skipped — one endpoint can fan out to different destinations based on content.

6. The Delivery Engine

The delivery worker runs in a tight loop: consume from transformed-events, acquire a global asyncio.Semaphore(50), and dispatch. Inside _deliver_with_retry, three resilience patterns run per destination URL.

Circuit Breaker

Three states stored in Redis:

This prevents cascading failures when a downstream service is degraded.

Sliding Window Rate Limiter

Redis sorted sets track request timestamps per destination:

def allow_request(url):
    now = now_ms()
    window_start = now - 60_000
    with redis.pipeline() as pipe:
        pipe.zremrangebyscore(key, 0, window_start)
        pipe.zcard(key)
        pipe.zadd(key, {uuid4(): now})
        results = await pipe.execute()
    return results[1] < max_rpm

The pipeline ensures atomicity. A 60-second TTL on the key prevents unbounded memory growth.

Exponential Backoff with Jitter

delay = (backoff_ms * (2 ** attempt)) / 1000
delay += random.uniform(0, delay * 0.5)  # jitter
await asyncio.sleep(delay)

Jitter prevents the thundering herd problem when many events fail against the same destination simultaneously.

Dead Letter Queue

After max_retries consecutive failures, the event is published to a dead-letter topic. The event body and full delivery history remain in PostgreSQL. Operators can:

7. SSRF Protection: Defending the Internal Network

The relay makes HTTP requests to user-configured URLs. Without protection, an attacker could point a route at http://localhost:5432 and exfiltrate the database.

url_security.py enforces:

  1. Scheme whitelist — only http/https, with HTTPS required in production
  2. Credential rejection — URLs with embedded user:password@ are rejected
  3. Hostname blocklist — localhost, localhost.localdomain, and raw private IPs are blocked
  4. DNS resolution check — at request time, the hostname is resolved and each resolved address is checked against private/reserved IP ranges
def _is_blocked_ip(address: str) → bool:
    ip = ipaddress.ip_address(address)
    return (
        ip.is_private or ip.is_loopback
        or ip.is_link_local or ip.is_multicast
        or ip.is_reserved or ip.is_unspecified
    )

This catches both explicit SSRF attempts and DNS rebinding attacks (where a hostname resolves to a public address at config time but a private address at request time).

8. Multi-Tenancy & RBAC

The service supports workspaces, users, and role-based access:

Role Capabilities
owner Full control, can delete workspace
admin Create/update/delete endpoints and routes, manage members
member Create endpoints and routes within workspace
viewer Read-only access to events, routes, delivery attempts

Registration auto-creates a personal workspace. Invitations add users with a specified role.

Authentication uses JWT with:

9. Why Not Microservices?

Microservices add service discovery, inter-service auth, distributed tracing — you'd spend months on infra before writing any business logic.

Instead, the codebase has clear package boundaries: app/gateway, app/transform, app/delivery, workers/. Each worker entry point is a standalone python -m workers.transform_worker command. Scaling means running the same Docker image with different commands:

docker run relay delivery-worker   # instance 1
docker run relay delivery-worker   # instance 2

Kafka consumer groups handle partition rebalancing. No service mesh, no sidecars, no distributed tracing until traffic justifies it.

10. Operational Considerations

Observability

Comes from PostgreSQL, not a metrics pipeline. The DeliveryAttempt table is an immutable audit log — every HTTP request is recorded with status, duration, and error message. The DLQ view is a SQL join between events and delivery_attempts, not a separate Kafka consumer. The React dashboard (Vite + Tailwind) queries this directly via the REST API.

Idempotent Migrations

Run on every startup:

await conn.execute(
    text("ALTER TABLE events ADD COLUMN IF NOT EXISTS "
         "status VARCHAR(32) NOT NULL DEFAULT 'pending';")
)

No Alembic. For a Docker-first deployment where each deploy is a fresh container, CREATE TABLE IF NOT EXISTS + ADD COLUMN IF NOT EXISTS is simpler and safer.

11. Lessons Learned

  1. Decouple early. The Kafka layer took one afternoon to add but saved weeks of refactoring later. The gateway never needed to know about routes, transforms, or retries.
  2. Optional infrastructure. Every external dependency (Kafka, Redis) can be missing at startup. The service falls back gracefully and logs clearly.
  3. PostgreSQL as source of truth, Kafka as accelerator. Events always exist in the database. Kafka is just a notification bus. If the bus breaks, the events are still there.
  4. Redis for coordination. Circuit breaker state, rate limit windows, idempotency keys, auth rate limits, token blocklists — all shared state lives in Redis. It's the one service that must always be available.
  5. SSRF is a first-class concern. Don't add HTTP delivery without URL validation. The DNS resolution check catches what string parsing misses.
  6. Event status tracking prevents silent data loss. Without pending → queued/failed, you can't tell whether an event was successfully handed off to Kafka. A status column is a single string column that saves hours of debugging.
  7. Fire-and-forget with an escape hatch. The gateway publishes to Kafka asynchronously so the 202 Accepted response isn't blocked. But the status column and replay API ensure no event is truly lost.

12. What to Improve and Implement Later

The current architecture handles the core webhook relay use case well, but several enhancements would make it production-ready at scale:

Sandboxed JavaScript Transformations

The current template engine supports dot-path access and division. Real-world transformation needs are more complex — string formatting, date manipulation, conditional logic, array mapping. A sandboxed JavaScript runtime (via quickjs or PyMiniRacer) would let users write arbitrary transform functions with memory and CPU limits, without opening an eval() security hole.

Prometheus Metrics & Structured Observability

Currently, observability comes from PostgreSQL audit logs. This is great for debugging individual events but terrible for dashboards. Exposing /metrics with ingestion rates, delivery latency percentiles (p50/p95/p99), worker queue depth, circuit breaker state counts, and Kafka consumer lag would enable real-time alerting via Grafana. Each worker should also export health and readiness probes so orchestrators (K8s, Nomad) can manage them properly.

DLQ Management Dashboard UI

The DLQ API exists (list, discard, restore, replay) but there's no dedicated dashboard view. A dead-letter queue page with search, batch replay, payload inspection before replay, and discard-all actions would drastically improve the operator experience. The current event viewer in the dashboard is read-only.

Webhook Retry Schedule (retry_at)

The retry_at column on events is defined but not yet used by a scheduler. A periodic job (or a separate retry worker) could query WHERE status = 'failed' AND retry_at < now() and automatically re-publish those events to Kafka. This turns failed from a terminal state into a transient one.

Schema Validation for Incoming Webhooks

Endpoints currently accept any valid JSON payload. Adding an optional JSON Schema (request_body_schema) per endpoint would let operators reject malformed payloads early with a 422, rather than discovering the problem when the delivery worker fails to transform the body.

Per-Endpoint Rate Limiting

The current rate limiter operates per destination URL (one route). A webhook source like Stripe can fire thousands of events per second — without a per-endpoint rate limit, a burst could overwhelm the delivery worker's global semaphore and starve other endpoints.

Batched Delivery

Some destinations support receiving events in batches (e.g., POST a JSON array of events instead of individual POSTs). A batch_window_ms route option would buffer events per destination and flush them as a single HTTP request, dramatically reducing connection overhead for high-volume integrations.

Error Notifications (Webhook Failures)

When an event lands in the DLQ, the only way to know is by polling the API. Integrating with Slack, PagerDuty, or email for DLQ alerts would give operators proactive visibility into delivery failures without staring at a dashboard.

Distributed Tracing

An event crosses three processes (gateway → transform → delivery). Without trace IDs propagated via Kafka headers, correlating a single event across all three is manual SQL joins. Adding OpenTelemetry instrumentation with trace_id in Kafka message headers would make debugging latency spikes and failures much faster.


The full source code is available on GitHub. docker compose up boots the entire stack in under 30 seconds.