Contents
1. The Problem 2. Architecture Overview 3. Ingestion Gateway 4. Kafka as the Decoupling Layer 5. The Transform Pipeline 6. The Delivery Engine 7. SSRF Protection 8. Multi-Tenancy & RBAC 9. Why Not Microservices? 10. Operational Considerations 11. Lessons Learned1. 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
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?
- Consumer groups let us scale workers horizontally. Spin up 3 delivery workers and Kafka partitions the load automatically.
- At-least-once delivery with auto-commit means a worker crash loses at most the current batch.
- KRaft mode eliminates ZooKeeper — a single Docker Compose service.
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:
- CLOSED — normal operation
- OPEN — after
Nconsecutive failures (circuit_breaker_threshold, default 10); all requests rejected forcooldown_sseconds - HALF_OPEN — after cooldown, one probe request is let through; success closes, failure re-opens
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:
- Replay via
POST /api/events/{id}/replay - Discard (soft-delete, hides from DLQ view)
- Restore a previously discarded event
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:
- Scheme whitelist — only
http/https, with HTTPS required in production - Credential rejection — URLs with embedded
user:password@are rejected - Hostname blocklist —
localhost,localhost.localdomain, and raw private IPs are blocked - 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:
- Peppered password hashing —
HMAC-SHA256(pepper, password)→bcrypt— the pepper is server-side only, so a database leak alone can't crack passwords - Token revocation via Redis blocklist
(
revoked:jti:{id}) — stateless JWTs with an escape hatch for immediate invalidation - Auth rate limiting — per-IP and per-email, 10 requests/minute on login/register
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
- 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.
- Optional infrastructure. Every external dependency (Kafka, Redis) can be missing at startup. The service falls back gracefully and logs clearly.
- 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.
- 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.
- SSRF is a first-class concern. Don't add HTTP delivery without URL validation. The DNS resolution check catches what string parsing misses.
- 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. - Fire-and-forget with an escape hatch. The gateway
publishes to Kafka asynchronously so the
202 Acceptedresponse 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.