Webhook Project
DodaTech
6 min read
title: "Webhook Mini Project: Complete Webhook System" description: "Build a complete webhook system from scratch including a provider, consumer, dead letter queue, monitoring dashboard, and security measures." weight: 30 date: 2026-06-28 lastmod: 2026-06-28 tags: ["apis", "webhooks"]
This project brings together everything you have learned about webhooks. You will build a complete webhook system with a provider that generates events, a consumer that receives them, a delivery engine with retry and DLQ, monitoring metrics, and comprehensive security. The project simulates a real-world e-commerce event notification system.
## What You'll Learn
- Integrate provider, consumer, and monitoring components into a unified system
- Implement end-to-end webhook delivery with retry and dead letter handling
- Add security at every layer of the webhook pipeline
- Build operational dashboards for webhook observability
## Why It Matters
A complete webhook system is more than the sum of its parts. The interplay between provider, consumer, security, and monitoring creates a production-grade integration platform. This project demonstrates how all the patterns work together and prepares you to build real webhook infrastructure.
## Real-World Use
- E-commerce platforms combine order webhooks (provider) with inventory sync (consumer)
- Payment gateways connect transaction events (provider) with accounting systems (consumer)
- SaaS platforms use webhooks to sync data across CRM, billing, and support tools
- DevOps pipelines link source control events (provider) with CI/CD systems (consumer)
## System Architecture
```mermaid
graph TD
subgraph "Provider"
A[Event Generator] --> B[Delivery Engine]
B --> C[Retry Queue]
C --> D[Dead Letter Queue]
end
subgraph "Consumer"
E[Webhook Endpoint]
F[Event Processor]
G[Audit Database]
end
subgraph "Security"
H[Signature Signing]
I[Signature Verification]
J[IP Allowlist]
end
subgraph "Monitoring"
K[Metrics Export]
L[Structured Logging]
M[Dashboard]
end
A --> H
H --> B
B --> E
E --> I
I --> J
J --> F
F --> G
B --> K
E --> K
K --> M
B --> L
E --> L
Project Requirements
Build a webhook system with the following components:
Provider Service (port 5000)
- Event generator that creates order events (created, updated, cancelled, refunded)
- Subscription management API (register, list, delete subscribers)
- Delivery engine with configurable retry (max 5 retries, exponential backoff)
- Per-consumer rate limiting (configurable events per second)
- HMAC-SHA256 signature signing with timestamp
- Dead letter queue for events that fail all retries
- Metrics endpoint at
/metrics(Prometheus format)
Consumer Service (port 5001)
- Webhook endpoint at
/webhookwith HMAC-SHA256 signature verification - Timestamp-based replay attack protection (5-minute tolerance)
- Payload validation and sanitization
- Idempotent event processing with deduplication
- Event storage in SQLite database
- Health check endpoint at
/health - Metrics endpoint at
/metrics(Prometheus format)
Monitoring
- Prometheus metrics for delivery count, success rate, latency, retries, and DLQ size
- Structured JSON logging for all events and deliveries
- Grafana dashboard (provisioned via JSON) showing:
- Delivery success rate over time
- Latency P50/P95/P99
- Queue depth and DLQ count
- Events per type distribution
- Top failing consumers
Implementation Plan
Step 1: Provider Core
import hashlib
import hmac
import json
import time
import uuid
from flask import Flask, request, jsonify
from prometheus_client import Counter, Histogram, Gauge, generate_latest
provider_app = Flask(__name__)
delivery_counter = Counter("webhook_delivery_total", "Total deliveries", ["consumer", "status"])
latency_histogram = Histogram("webhook_delivery_latency", "Delivery latency", ["consumer"])
dlq_gauge = Gauge("webhook_dlq_size", "Dead letter queue size")
SECRET = "provider_secret_key_2026"
def sign_payload(payload):
timestamp = int(time.time())
data = f"{timestamp}.{json.dumps(payload, sort_keys=True)}".encode()
signature = hmac.new(SECRET.encode(), data, hashlib.sha256).hexdigest()
return timestamp, f"t={timestamp},v1={signature}"
subscriptions = []
delivery_queue = []
dlq = []
@provider_app.route("/api/subscribe", methods=["POST"])
def subscribe():
data = request.get_json()
sub = {
"id": str(uuid.uuid4()),
"url": data["url"],
"events": data.get("events", ["*"]),
"rate_limit": data.get("rate_limit", 10),
"created_at": time.time()
}
subscriptions.append(sub)
return jsonify(sub), 201
@provider_app.route("/api/events", methods=["POST"])
def create_event():
event = request.get_json()
event["id"] = f"evt_{uuid.uuid4().hex[:12]}"
event["created_at"] = time.time()
for sub in subscriptions:
if "*" in sub["events"] or event["type"] in sub["events"]:
delivery_queue.append({"sub": sub, "event": event})
return jsonify({"status": "queued", "event_id": event["id"]}), 202
@provider_app.route("/metrics")
def metrics():
return generate_latest(), 200, {"Content-Type": "text/plain"}
if __name__ == "__main__":
provider_app.run(port=5000)
Step 2: Consumer Core
import hashlib
import hmac
import json
import sqlite3
import time
from flask import Flask, request, jsonify
consumer_app = Flask(__name__)
SECRET = "provider_secret_key_2026"
DEDUP_DB = sqlite3.connect(":memory:", check_same_thread=False)
DEDUP_DB.execute("CREATE TABLE processed_events (event_id TEXT PRIMARY KEY, processed_at REAL)")
DEDUP_DB.execute("CREATE TABLE received_events (id INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT, type TEXT, payload TEXT, received_at REAL)")
def verify_signature(payload, signature_header, tolerance=300):
try:
parts = signature_header.split(",")
ts = None
sig = None
for part in parts:
if part.startswith("t="):
ts = int(part[2:])
elif part.startswith("v1="):
sig = part[3:]
if ts is None or sig is None:
return False
if abs(time.time() - ts) > tolerance:
return False
data = f"{ts}.{json.dumps(payload, sort_keys=True)}".encode()
expected = hmac.new(SECRET.encode(), data, hashlib.sha256).hexdigest()
return hmac.compare_digest(expected, sig)
except Exception:
return False
@consumer_app.route("/webhook", methods=["POST"])
def webhook():
signature = request.headers.get("X-Signature", "")
payload = request.get_json()
if not payload:
return jsonify({"error": "invalid payload"}), 400
if not verify_signature(payload, signature):
return jsonify({"error": "invalid signature"}), 401
event_id = payload.get("id")
if not event_id:
return jsonify({"error": "missing event id"}), 400
cur = DEDUP_DB.execute("SELECT 1 FROM processed_events WHERE event_id = ?", (event_id,))
if cur.fetchone():
return jsonify({"status": "duplicate"}), 200
DEDUP_DB.execute("INSERT INTO processed_events VALUES (?, ?)", (event_id, time.time()))
DEDUP_DB.execute("INSERT INTO received_events (event_id, type, payload, received_at) VALUES (?, ?, ?, ?)",
(event_id, payload["type"], json.dumps(payload), time.time()))
DEDUP_DB.commit()
return jsonify({"status": "processed"}), 200
@consumer_app.route("/health")
def health():
return jsonify({"status": "healthy", "events_processed": "ok"}), 200
if __name__ == "__main__":
consumer_app.run(port=5001)
Step 3: Delivery Engine with Retry and DLQ
import time
import requests
import threading
import json
class DeliveryEngine:
def __init__(self, provider):
self.provider = provider
self.max_retries = 5
self.base_delay = 2
self.running = True
def process_queue(self):
while self.running:
for task in self.provider.delivery_queue[:]:
event = task["event"]
sub = task["sub"]
payload = {
"id": event["id"],
"type": event["type"],
"data": event.get("data", {}),
"created_at": event["created_at"]
}
timestamp, signature = self.provider.sign_payload(payload)
start = time.time()
try:
resp = requests.post(
sub["url"],
json=payload,
headers={
"Content-Type": "application/json",
"X-Signature": signature,
"X-Event-ID": event["id"]
},
timeout=10
)
latency = time.time() - start
if resp.status_code // 100 == 2:
self.provider.delivery_queue.remove(task)
self.provider.delivery_counter.labels(consumer=sub["url"], status="success").inc()
self.provider.latency_histogram.labels(consumer=sub["url"]).observe(latency)
else:
task.setdefault("attempts", 0)
task["attempts"] += 1
if task["attempts"] >= self.max_retries:
self.provider.delivery_queue.remove(task)
self.provider.dlq.append(task)
self.provider.dlq_gauge.set(len(self.provider.dlq))
self.provider.delivery_counter.labels(consumer=sub["url"], status="dlq").inc()
else:
time.sleep(self.base_delay * (2 ** (task["attempts"] - 1)))
self.provider.delivery_counter.labels(consumer=sub["url"], status="retry").inc()
except Exception as e:
task.setdefault("attempts", 0)
task["attempts"] += 1
if task["attempts"] >= self.max_retries:
self.provider.delivery_queue.remove(task)
self.provider.dlq.append(task)
self.provider.dlq_gauge.set(len(self.provider.dlq))
else:
time.sleep(self.base_delay * (2 ** (task["attempts"] - 1)))
time.sleep(0.5)
Step 4: Test Script
import requests
import time
import json
BASE = "http://localhost:5000"
CONSUMER = "http://localhost:5001"
sub = requests.post(f"{BASE}/api/subscribe", json={
"url": f"{CONSUMER}/webhook",
"events": ["order.created", "order.updated"]
}).json()
print(f"Subscribed: {sub['id']}")
for i in range(5):
event = {
"type": "order.created",
"data": {"order_id": i, "amount": 100 + i, "currency": "USD"}
}
resp = requests.post(f"{BASE}/api/events", json=event)
print(f"Event {i}: {resp.json()}")
time.sleep(3)
health = requests.get(f"{CONSUMER}/health")
print(f"Consumer health: {health.json()}")
metrics = requests.get(f"{BASE}/metrics")
print(f"Provider metrics:\n{metrics.text[:500]}")
metrics = requests.get(f"{CONSUMER}/metrics")
print(f"Consumer metrics:\n{metrics.text[:500]}")
Step 5: Monitor and Verify
- Start both provider and consumer services
- Run the test script to generate events
- Check that events appear in the consumer's database
- Verify Prometheus metrics at
/metricsendpoints - Test security by sending unsigned requests (should return 401)
- Test replay protection by sending the same request twice (should return duplicate)
- Test DLQ by pointing to a non-existent consumer URL
Common Mistakes
- Not testing the end-to-end flow with the consumer offline to verify DLQ behavior
- Forgetting to configure CORS or firewall rules between provider and consumer
- Using different secrets for signing and verification between provider and consumer
- Not handling the case where the consumer returns a non-standard HTTP status code
- Overlooking the need for database connection pooling in production
Project Deliverables
- Provider service (provider.py) with event generation, delivery engine, DLQ, metrics, and subscription API
- Consumer service (consumer.py) with webhook endpoint, signature verification, dedup, and health check
- Delivery engine (engine.py) with retry logic and DLQ management
- Test script (test.py) to exercise the complete system
- Docker compose (docker-compose.yml) to run provider, consumer, and Prometheus + Grafana
- Grafana dashboard (dashboard.json) showing webhook delivery metrics
What's Next
Congratulations on completing the webhooks module. You are now ready to explore more API topics or apply these patterns to build production webhook systems.
Built by the developers of DodaTech
Doda Browser, DodaZIP & Durga Antivirus Pro