Reliable Webhook Delivery with a Job Queue
Outbound webhooks are the purest retry problem in background processing: you call endpoints you do not control, many of which will be slow, down, or misconfigured, and customers expect every event eventually. This guide builds a webhook delivery system on a job queue, as part of Retry Strategies & Backoff in Queue Fundamentals & Architecture.
Problem Statement
A SaaS platform sends webhooks for events such as invoice.paid to about 4,000 customer endpoints. The first implementation posted from a Celery task with three retries over one minute. Customers complained about missed events whenever their endpoint was down for more than a minute. When one large customer's endpoint started taking 30 seconds to respond, it tied up most worker slots and delayed webhooks for every other customer by hours. And some customers received duplicate events after worker restarts, with no way to tell them apart. You want at-least-once delivery with retries spanning a day, isolation between endpoints, signatures and idempotency ids customers can rely on, automatic disabling of dead endpoints, and a delivery log customers and support can see.
Prerequisites
- A job framework with delayed retries (Celery, Sidekiq, BullMQ, or a database-backed queue).
- A table for webhook endpoints (URL, secret, status) and one for delivery attempts.
- An outbox or transactional enqueue so events are not lost between the business write and the queue.
- An HTTP client with strict connect and read timeouts.
Step 1 — Record Events and Deliveries Durably
Webhooks start as business events. Write the event and one delivery row per subscribed endpoint in the same transaction as the business change, then enqueue delivery jobs from those rows. The delivery row, not the queue message, is the source of truth.
CREATE TABLE webhook_events (
id uuid PRIMARY KEY, -- sent to customers as the idempotency id
type text NOT NULL,
payload jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE TABLE webhook_deliveries (
id bigserial PRIMARY KEY,
event_id uuid REFERENCES webhook_events(id),
endpoint_id bigint NOT NULL,
status text NOT NULL DEFAULT 'pending', -- pending | delivered | failed | abandoned
attempts int NOT NULL DEFAULT 0,
next_try_at timestamptz NOT NULL DEFAULT now(),
last_status int, last_error text,
UNIQUE (event_id, endpoint_id)
);
Because deliveries are rows, a lost queue message is recoverable: a sweeper re-enqueues pending deliveries whose next_try_at has passed. This is the transactional outbox pattern applied to webhooks.
Step 2 — Sign Payloads and Send a Stable Event Id
Customers need to verify that a webhook came from you and to deduplicate retries. Send the event id as a header that stays the same across every attempt, and an HMAC signature over the timestamp and body.
import hmac, hashlib, json, time, httpx
def deliver(delivery_id: int) -> None:
d = load_delivery(delivery_id)
if d.status != "pending":
return # already delivered/abandoned
body = json.dumps({"id": str(d.event.id), "type": d.event.type,
"created": d.event.created_at.isoformat(), "data": d.event.payload},
separators=(",", ":")).encode()
ts = str(int(time.time()))
sig = hmac.new(d.endpoint.secret.encode(), f"{ts}.".encode() + body, hashlib.sha256).hexdigest()
headers = {"Content-Type": "application/json",
"Webhook-Id": str(d.event.id), # same on every retry: customers dedupe on it
"Webhook-Timestamp": ts,
"Webhook-Signature": f"v1={sig}",
"User-Agent": "ExampleWebhooks/1.0"}
try:
resp = client.post(d.endpoint.url, content=body, headers=headers,
timeout=httpx.Timeout(10.0, connect=3.0))
outcome = ("delivered", resp.status_code) if 200 <= resp.status_code < 300 else ("retry", resp.status_code)
except httpx.HTTPError as exc:
outcome = ("retry", None); d.last_error = type(exc).__name__
record_attempt(d, *outcome)
Including the timestamp in the signed content lets customers reject replays older than a few minutes. The Webhook-Id header is what makes at-least-once delivery acceptable to customers: a duplicate after a worker crash carries the same id, and their handler can ignore it — the idempotency pattern from preventing duplicate job execution with idempotency, exported to customers.
Step 3 — Retry Over Hours with Full Jitter
Customer endpoints go down for deploys, incidents, and weekends. Retries must span long enough to cover real outages — a day or more — with delays that grow and are jittered so a customer's recovery is not greeted by a storm.
SCHEDULE = [30, 120, 600, 1800, 3600, 7200, 14400, 28800, 43200, 86400] # seconds; ~2 days total
def record_attempt(d, outcome: str, status: int | None) -> None:
d.attempts += 1
d.last_status = status
if outcome == "delivered":
d.status = "delivered"
elif status in (400, 401, 403, 404, 410) and d.attempts >= 3:
d.status = "failed" # consistent client errors: stop early
elif d.attempts > len(SCHEDULE):
d.status = "failed"
else:
base = SCHEDULE[d.attempts - 1]
d.next_try_at = now() + timedelta(seconds=random.uniform(base * 0.5, base))
deliver_webhook.apply_async(args=[d.id], eta=d.next_try_at)
save(d)
update_endpoint_health(d.endpoint_id, outcome == "delivered")
The schedule is data, not code, which makes it easy to document for customers ("we retry for up to 48 hours"). Treating repeated 4xx responses as failures avoids retrying for two days against an endpoint that answers every request with 404. The jitter follows exponential backoff with jitter in Celery.
Step 4 — Isolate Endpoints So One Slow Customer Cannot Block Others
The hours-long delays in the problem statement came from one slow endpoint occupying worker slots. Two controls fix it: a short timeout (10 seconds is generous for a webhook), and a per-endpoint concurrency limit so no endpoint can hold more than a few slots at once.
MAX_INFLIGHT_PER_ENDPOINT = 4
@app.task(bind=True, acks_late=True)
def deliver_webhook(self, delivery_id: int):
d = load_delivery(delivery_id)
slot = acquire_semaphore(f"wh:inflight:{d.endpoint_id}", MAX_INFLIGHT_PER_ENDPOINT, ttl=30)
if slot is None:
# this endpoint is saturated: try again shortly, no attempt consumed
raise self.retry(countdown=random.uniform(5, 15), max_retries=None)
try:
deliver(delivery_id)
finally:
release_semaphore(f"wh:inflight:{d.endpoint_id}", slot)
With 4 slots per endpoint and 10-second timeouts, a completely hung endpoint consumes at most 4 slots and delays only its own deliveries. Everyone else's webhooks flow on the remaining capacity. The fairness principle is the same as in preventing tenant starvation with weighted queues.
Step 5 — Disable Dead Endpoints and Notify Owners
An endpoint that has failed every attempt for days is not coming back on its own. Track per-endpoint health; after a sustained failure period, disable it, stop enqueueing new deliveries, and email the account owner.
def update_endpoint_health(endpoint_id: int, ok: bool) -> None:
if ok:
db.execute("UPDATE webhook_endpoints SET failing_since = NULL WHERE id = %s", (endpoint_id,))
return
db.execute("""UPDATE webhook_endpoints SET failing_since = COALESCE(failing_since, now())
WHERE id = %s""", (endpoint_id,))
ep = load_endpoint(endpoint_id)
if ep.failing_since and now() - ep.failing_since > timedelta(days=3) and ep.status == "active":
db.execute("UPDATE webhook_endpoints SET status = 'disabled' WHERE id = %s", (endpoint_id,))
notify_owner(ep.account_id, "Webhook endpoint disabled after 3 days of failures", ep.url)
abandon_pending(endpoint_id) # mark pending deliveries abandoned
Disabled endpoints keep their delivery history. When the customer fixes and re-enables the endpoint, offer a "resend events since" action that creates new deliveries for abandoned events — replay on request, not automatically.
Step 6 — Expose a Delivery Log
Most webhook support tickets are "did you send it?". A per-endpoint delivery log answering that question — event id, attempts, response codes, timings — removes the ticket entirely.
-- The query behind the customer-facing log
SELECT e.id AS event_id, e.type, d.status, d.attempts, d.last_status, d.next_try_at, d.last_error
FROM webhook_deliveries d JOIN webhook_events e ON e.id = d.event_id
WHERE d.endpoint_id = $1 ORDER BY e.created_at DESC LIMIT 100;
Keep attempt-level detail (response code, latency, truncated response body) in a separate attempts table with shorter retention, and show it on demand.
Verification
def test_slow_endpoint_does_not_block_others(fake_endpoints, run_workers):
fake_endpoints.set_latency("slow.example", 30)
enqueue_events(endpoints=["slow.example"] + [f"ok{i}.example" for i in range(50)], count=200)
run_workers(slots=20, seconds=30)
assert fake_endpoints.delivered_fraction(exclude="slow.example") > 0.99
assert fake_endpoints.max_concurrent("slow.example") <= 4
In production, track delivery success rate per attempt number and the age of the oldest pending delivery; a rising first-attempt failure rate across many endpoints usually means a problem on your side (signing, outbound network) rather than theirs.
Gotchas & Edge Cases
Following redirects. Redirects can send payloads to unexpected hosts. Do not follow them by default; treat 3xx as a failure the customer must fix.
Server-side request forgery. Customer-supplied URLs can point at internal addresses. Resolve and block private ranges before connecting, and send from an egress path without access to internal services.
Ordering. Webhooks may arrive out of order under retries. Include a sequence or created timestamp and document that consumers should not rely on arrival order.
Payload size and secrets. Keep payloads small (send ids and let customers fetch details), and rotate signing secrets with an overlap period where both are accepted.
FAQ
Should one job deliver to all endpoints for an event? No. One job per (event, endpoint) isolates failures and retries; a combined job would retry successful endpoints when one fails.
How long should retries last? At least 24 hours, commonly up to three days, so customers can recover from weekend outages. Publish the schedule in your documentation.
How do customers verify the signature?
Publish a short snippet per language: recompute HMAC-SHA256(secret, timestamp + "." + raw_body), compare it to the header with a constant-time comparison, and reject timestamps older than five minutes. The most common customer bug is verifying against a re-serialized JSON body instead of the raw bytes received, so call that out explicitly in the documentation.
Is a managed webhook service worth it? For many teams, yes: services like Svix or Hookdeck implement signing, retries, isolation, and logs. The design here is what they do internally.
Related
- Retry Strategies & Backoff — the retry principles applied here.
- Google Cloud Tasks for HTTP Workers — managed per-queue rate limits for outbound calls.
- Circuit Breakers for Worker Dependencies — when your own egress is the failing dependency.
- Dead-Letter Queues & Poison Messages — what "failed" deliveries become.