The Fragility of Naive Webhook Dispatchers
In modern SaaS architectures and API platforms, webhooks are the lifeblood of real-time event integration: payment confirmations, invoice settlements, shipment status updates, and user actions. Despite their criticality, many engineering teams implement webhook dispatching through naive patterns:
- In-Band HTTP Requests: Firing an outbound HTTP request directly inside a web view. If the subscriber's endpoint is slow or offline, your web worker thread hangs, exhausting connection pools.
- Direct Celery/Queue Enqueuing: Calling
send_webhook_task.delay(event_data)inside a database transaction. If the transaction fails to commit or rolls back due to a constraint violation, the background worker still fires the webhook, dispatching a "ghost event" for data that was never saved. Conversely, if the message broker (RabbitMQ/Redis) crashes between the database commit and the enqueue call, the event is permanently lost.
The Transactional Outbox Pattern
To guarantee At-Least-Once Delivery without ghost events or silent loss, enterprise architectures utilize the Transactional Outbox Pattern. Instead of sending messages across the network during business logic execution, the application persists the outgoing webhook payload directly into a dedicated database table (webhook_outbox) as part of the exact same ACID database transaction that mutates business records.
If the transaction succeeds, the outbox record is guaranteed to exist on disk. If the transaction rolls back, the outbox record rolls back with it. A decoupled, asynchronous background dispatcher sweeps the outbox table, dispatches the HTTP payloads with retry backoffs, and marks successful deliveries.
Database Schema and Outbox Producer
Below is the PostgreSQL schema and Django model implementation utilizing row-level locking for multi-worker concurrency:
# models.py
from django.db import models
import uuid
class WebhookEndpoint(models.Model):
url = models.URLField(max_length=500)
secret_key = models.CharField(max_length=64) # Used for HMAC signing
is_active = models.BooleanField(default=True)
class WebhookOutboxEvent(models.Model):
id = models.UUIDField(primary_key=True, default=uuid.uuid4, editable=False)
endpoint = models.ForeignKey(WebhookEndpoint, on_delete=models.CASCADE)
event_type = models.CharField(max_length=100) # e.g. "payment.succeeded"
payload = models.JSONField()
status = models.CharField(
max_length=20,
choices=[('PENDING', 'Pending'), ('PROCESSING', 'Processing'), ('DELIVERED', 'Delivered'), ('FAILED', 'Failed')],
default='PENDING',
db_index=True
)
attempt_count = models.PositiveIntegerField(default=0)
next_retry_at = models.DateTimeField(auto_now_add=True, db_index=True)
created_at = models.DateTimeField(auto_now_add=True)
class Meta:
indexes = [
models.Index(fields=['status', 'next_retry_at']),
]
High-Throughput Dispatcher with SKIP LOCKED and HMAC Signing
To scale webhook dispatching across multiple concurrent workers without duplicate sends, workers claim pending events using PostgreSQL's FOR UPDATE SKIP LOCKED:
# dispatcher.py
import hmac
import hashlib
import time
import requests
from django.db import transaction
from django.utils import timezone
from .models import WebhookOutboxEvent
def process_webhook_batch(batch_size=50):
now = timezone.now()
with transaction.atomic():
# Atomically select and lock unclaimed pending events
events = list(
WebhookOutboxEvent.objects.select_for_update(skip_locked=True)
.filter(status='PENDING', next_retry_at__lte=now)
.select_related('endpoint')[:batch_size]
)
for event in events:
event.status = 'PROCESSING'
event.save(update_fields=['status'])
# Dispatch outside the database transaction
for event in events:
dispatch_single_webhook(event)
def dispatch_single_webhook(event: WebhookOutboxEvent):
payload_bytes = json.dumps(event.payload, separators=(',', ':')).encode('utf-8')
timestamp = str(int(time.time()))
signature_payload = f"t={timestamp}.".encode('utf-8') + payload_bytes
# Calculate HMAC-SHA256 signature to defend against tampering and replay attacks
signature = hmac.new(
event.endpoint.secret_key.encode('utf-8'),
signature_payload,
hashlib.sha256
).hexdigest()
headers = {
'Content-Type': 'application/json',
'X-DevManue-Signature': f"t={timestamp},v1={signature}",
'User-Agent': 'devManue-Webhook-Engine/1.0'
}
try:
response = requests.post(event.endpoint.url, data=payload_bytes, headers=headers, timeout=5)
if 200 <= response.status_code < 300:
event.status = 'DELIVERED'
event.save(update_fields=['status'])
return
except requests.RequestException:
pass # Fall through to exponential retry backoff
# Calculate exponential backoff with jitter: 2^attempt * 15 seconds
event.attempt_count += 1
if event.attempt_count >= 5:
event.status = 'FAILED'
else:
event.status = 'PENDING'
backoff_seconds = (2 ** event.attempt_count) * 15
event.next_retry_at = timezone.now() + timezone.timedelta(seconds=backoff_seconds)
event.save(update_fields=['status', 'attempt_count', 'next_retry_at'])
By decoupling generation from dispatch and using cryptographic signatures, the webhook subsystem operates with mathematical reliability. For high-throughput API integrations, review our High-Throughput APIs & Distributed Django Services.