Prompt
I'm building an integration layer for webhook events with retries and filtering — what architecture makes sense?
Latest observation
A solid architecture for a webhook integration layer usually looks like an event ingestion + durable queue + worker pipeline + retry scheduler setup, with filtering applied as early as possible.
Recommended architecture
1) Webhook receiver / ingress service
This is the public HTTP endpoint that receives events from upstream providers.
Responsibilities:
- Verify authenticity:
- HMAC signatures
- API keys / mTLS / JWT, depending on provider
- Deduplicate:
- Use event ID + provider ID
- Store a processed/seen marker
- Acknowledge fast:
- Return
2xxquickly once the payload is safely persisted
- Return
- Normalize:
- Convert provider-specific payloads into a common internal envelope
Best practice: keep this layer thin and fast.
2) Durable event store or message queue
Do not process webhooks inline in the HTTP request.
Use a durable buffer such as:
- Kafka / Redpanda
- SQS / SNS
- RabbitMQ
- NATS JetStream
- A database-backed queue if scale is modest
Store:
- raw payload
- normalized event
- headers
- metadata like delivery attempt count, tenant, source, received timestamp, correlation ID
This gives you:
- retry safety
- backpressure handling
- replay capability
- auditability
3) Filtering / routing layer
Apply filters before expensive work.
Common filter types:
- event type
- tenant/customer
- source provider
- payload fields
- environment
- custom rules
Implementation options:
- simple rule engine
- declarative configs
- message topic routing
- predicate-based workers
Pattern:
- Ingest all events into the queue
- Evaluate filters in a routing worker
- Dispatch only matching events to downstream handlers
If filters are static and simple, route by topic/partition key. If filters are dynamic or customer-defined, use a policy/rules service.
4) Worker processing layer
Workers consume filtered events and perform the actual integration work:
- transform payload
- call internal services
- push to destination APIs
- update state in DB
- emit metrics/logs
Make workers:
- idempotent
- stateless
- horizontally scalable
Idempotency is essential because retries and redeliveries are normal.
5) Retry management
Use a retry policy that distinguishes between:
- transient failures: timeouts, 429s, 5xx
- permanent failures: invalid schema, auth failure, unsupported event type
Recommended retry approach:
- exponential backoff with jitter
- max retry count or max age
- dead-letter queue after exhaustion
- separate handling for poison messages
Example schedule:
- immediate retry once
- then 1m, 5m, 15m, 1h, 6h
- then DLQ
Important:
- retry at the worker/queue level, not in the ingress request
- preserve attempt history
6) Dead-letter queue / failure store
When processing fails repeatedly, move the event aside.
Store:
- original payload
- failure reason
- attempt history
- stack trace / error class
- timestamps
This enables:
- manual replay
- debugging
- observability
- customer support workflows
7) Replay / reprocessing tooling
You’ll want admin tools to:
- replay from DLQ
- replay by time range
- replay by tenant/event type
- reprocess after filter changes or code fixes
Make replay idempotent too.
Suggested data flow
A practical flow:
- Provider sends webhook
- Receiver validates signature
- Receiver stores raw event and enqueues normalized event
- Routing worker evaluates filters
- Matching events go to handler queues
- Workers process events
- Failures go through retry schedule
- Exhausted failures move to DLQ
- Replay tools can re-inject from DLQ
Key design principles
Idempotency
Assume duplicate delivery. Use:
- event IDs
- dedupe keys
- idempotency keys for outbound API calls
At-least-once delivery
Most systems will be at-least-once, not exactly-once. Design handlers accordingly.
Separation of concerns
Keep these separate:
- receipt
- persistence
- filtering
- processing
- retry orchestration
- observability
Fast ACKs
Respond quickly to the webhook sender to avoid provider retries.
Backpressure
If downstream is slow:
- queue absorbs load
- workers scale up
- apply rate limits per tenant/provider if needed
Filtering architecture choices
Option A: Filter in code
Good for:
- low complexity
- few event types
- centrally managed logic
Tradeoff:
- harder to customize per tenant
Option B: Rule engine
Good for:
- customer-configurable filters
- dynamic policies
- complex conditions
Tradeoff:
- more operational complexity
Option C: Queue/topic routing
Good for:
- very high throughput
- clear event categories
Tradeoff:
- less flexible for arbitrary rules
A simple reference architecture
Webhook Provider
|
v
Ingress API -- validate/signature/dedupe --> Raw Event Store
|
v
Queue / Stream
|
v
Router / Filter Evaluator
|
+--> Handler Queue A --> Worker A
+--> Handler Queue B --> Worker B
+--> ...
|
v
Retry Queue / Scheduler
|
v
Dead Letter Queue
|
v
Replay/Admin Tools
Practical implementation tips
- Use a DB transaction or outbox pattern if event persistence and queue publication must be atomic
- Include a schema version in the event envelope
- Track processing state with timestamps and attempt counts
- Add per-tenant rate limiting
- Log structured JSON with correlation IDs
- Expose metrics:
- ingest rate
- filter match rate
- success/failure rates
- retry count distribution
- queue lag
- DLQ volume
- Keep payloads immutable; store transformed copies separately
If you’re starting small
A good MVP architecture is:
- HTTP webhook receiver
- PostgreSQL table for raw events + processing status
- background workers
- retry state in DB
- dead-letter table
- simple filter rules in code
This is often enough before introducing Kafka or a dedicated queue.
Rule of thumb
If your main concerns are:
- reliability
- retries
- visibility
- customer-specific filters
then the best pattern is:
accept fast, persist durably, filter asynchronously, process idempotently, retry with backoff, dead-letter on exhaustion, replay from storage.
If you want, I can also sketch:
- a minimal MVP architecture,
- a cloud-native version on AWS/GCP, or
- a database schema + worker design for this system.