Prompt
ChatGPT, help me build a webhook pipeline that filters, transforms, and routes events to different services.
Latest observation
Absolutely — a webhook pipeline usually has 4 parts:
- Receive webhook events
- Validate + filter them
- Transform/enrich the payload
- Route to one or more downstream services
Here’s a practical blueprint you can use.
1) High-level architecture
Webhook Source → Ingress API → Filter/Transform Layer → Router → Destinations
Example flow:
- A provider sends
POST /webhooks/events - Your service:
- verifies signature
- checks event type
- drops unwanted events
- normalizes the payload
- routes it to:
- Slack for alerts
- CRM for customer events
- Queue for async processing
- Analytics sink for tracking
2) Core design principles
A. Fast acknowledgement
Webhook senders often retry if you respond slowly.
Return 2xx quickly after:
- validation
- basic filtering
- persisting the event or queueing it
Then process downstream asynchronously.
B. Idempotency
Webhooks are frequently duplicated.
Store an event_id and ignore repeats.
C. Separation of concerns
Keep these separate:
- ingestion
- filtering rules
- transformation
- routing
- delivery retries
D. Observability
Log:
- event id
- source
- destination
- result
- latency
- retry count
3) Suggested pipeline stages
Stage 1: Ingress
- Accept JSON webhook
- Validate authentication/signature
- Normalize headers and metadata
Example checks:
- HMAC signature
- timestamp freshness
- source IP allowlist
- content type
Stage 2: Filter
Filter based on rules like:
- event type
- source
- tenant/customer
- severity
- payload fields
Examples:
- only process
invoice.paid - ignore test mode events
- drop events missing
customer_id
Stage 3: Transform
Map source-specific payloads into a standard internal event format.
Example transformation:
user.email→customer.emailcreated_at→ ISO timestamp- enrich with internal tenant config or lookup data
Stage 4: Route
Send the transformed event to one or many targets.
Routing patterns:
- fan-out: send to all matching destinations
- conditional routing: send based on rules
- priority routing: critical events first
- dead-letter routing: failed events to a DLQ
4) Example internal event format
{
"event_id": "evt_123",
"source": "stripe",
"type": "invoice.paid",
"tenant_id": "t_001",
"timestamp": "2026-10-05T12:00:00Z",
"data": {
"customer_id": "cus_456",
"amount": 4999,
"currency": "usd"
},
"metadata": {
"environment": "prod"
}
}
5) Routing rules example
You can define rules in config or a database:
[
{
"match": { "type": "invoice.paid" },
"destinations": ["billing-service", "slack-alerts"]
},
{
"match": { "type": "user.created", "tenant_id": "t_001" },
"destinations": ["crm-service"]
},
{
"match": { "severity": "high" },
"destinations": ["pagerduty"]
}
]
6) Simple implementation pattern
Option A: Single service + queue
Good for most cases.
Flow:
- Webhook API writes event to DB/queue
- Worker processes filter/transform/route
- Retry failures automatically
Suggested tools:
- API: FastAPI / Express / Go / Flask
- Queue: SQS / RabbitMQ / Kafka / Redis Streams
- Storage: Postgres
- Observability: OpenTelemetry + logs
7) Example pseudo-code
Ingress handler
def webhook_handler(request):
verify_signature(request)
event = request.json()
if is_duplicate(event["event_id"]):
return {"status": "ok", "duplicate": True}
save_raw_event(event)
enqueue(event["event_id"])
return {"status": "accepted"}
Worker
def process_event(event_id):
event = load_raw_event(event_id)
if not should_process(event):
mark_skipped(event_id)
return
transformed = transform_event(event)
destinations = route_event(transformed)
for destination in destinations:
send_to_destination(destination, transformed)
mark_done(event_id)
8) Example filter/transform/router logic
def should_process(event):
return (
event.get("type") in {"invoice.paid", "user.created"} and
event.get("metadata", {}).get("environment") == "prod"
)
def transform_event(event):
return {
"event_id": event["event_id"],
"type": event["type"],
"customer_id": event["data"].get("customer_id"),
"amount_usd": event["data"].get("amount", 0) / 100,
"timestamp": event["timestamp"]
}
def route_event(event):
destinations = []
if event["type"] == "invoice.paid":
destinations.append("billing-service")
destinations.append("slack-alerts")
if event["amount_usd"] > 1000:
destinations.append("finance-review")
return destinations
9) Reliability features you’ll want
- Retry with exponential backoff
- Dead-letter queue
- Timeouts per destination
- Circuit breaker for flaky services
- Rate limiting
- Schema validation
- Versioned payloads
- Replay capability for failed events
10) Security checklist
- Verify webhook signatures
- Rotate secrets
- Reject old timestamps
- Validate JSON schema
- Sanitize outbound payloads
- Limit payload size
- Store raw payloads securely
- Separate tenants logically
11) If you want rule-based routing
A good pattern is a small routing DSL or config like:
rules:
- name: invoice-paid
when:
type: invoice.paid
send_to:
- billing-service
- slack-alerts
- name: vip-customer
when:
type: customer.updated
customer_tier: vip
send_to:
- account-manager
Then your router evaluates rules in order.
12) Recommended minimal stack
If you want to build this quickly:
- API: FastAPI
- Queue: Redis Streams or SQS
- Worker: Celery / Dramatiq / custom consumer
- DB: Postgres
- Validation: Pydantic / JSON Schema
- Delivery: HTTP client with retries
13) A good MVP plan
MVP v1
- Receive webhook
- Validate signature
- Store raw event
- Queue processing job
- Filter by event type
- Transform into internal schema
- Route to 2 destinations
- Retry failures
MVP v2
- UI/config for routing rules
- Per-destination auth
- Replay failed events
- Metrics dashboard
If you want, I can next give you one of these:
- a concrete implementation in Python/FastAPI,
- a Node.js/Express version,
- a system design diagram, or
- a database schema + routing rules model.