High-Throughput Webhook Systems: Delivery Timelines and Split-Brain Prevention
Learn how to debug asynchronous webhook jobs with full-lifecycle timeline tracking, worker heartbeats, job leasing, and split-brain prevention in queue pipelines.

High-Throughput Webhook Systems: Delivery Timelines and Split-Brain Prevention
Webhooks are the standard way to integrate systems asynchronously. Scaling a delivery platform to thousands of requests per second exposes two infrastructure problems:
- Debugging the "black box": Most webhook logs record only
ReceivedandDelivered. When a transient failure or partition occurs, engineers can't see the intermediate states (encrypted, leased, attempted, retried, dead-lettered). - Duplicate processing after lease expiry: In large worker clusters, a stalled thread or latency spike can let a lease expire while the original worker is still alive. A second worker then claims the job, and the receiver may see the same event twice.
This guide covers a state-machine design for full-lifecycle visibility. It also covers lease-and-fencing patterns that keep your own state consistent, and the part fencing cannot solve: duplicates seen by the receiver.
1. Full-Lifecycle Timelines for Asynchronous Webhook Jobs
The limits of Received vs. Delivered logging
A basic setup writes to a central log sink (Elasticsearch, CloudWatch, Loki) at two points:
- When the payload reaches the ingest API (
Received). - When the receiver answers with an HTTP
2xxstatus (Delivered).
This works at low volume but breaks down at scale:
- Unaccounted queue latency: If an event sits for 30 seconds before a worker picks it up, two-point logging can't show whether the cause was queue lag, database lock contention, or worker starvation.
- Opaque retry behavior: After a
503, a retry is scheduled. Without state transitions, operators can't tell a job sleeping in backoff from one whose worker crashed. - Audit gaps: If payloads contain sensitive data, you may need to show that encryption happened before storage. Without structured state events, that means correlating logs across several services.
Retry windows at the sending side are long, so visibility matters. Stripe, for example, attempts to deliver a given event to your webhook endpoint for up to 3 days with an exponential back off in live mode, while test mode retries three times over a few hours. A job can stay in flight for days, and operators need to see where it stands. stripe
A state-machine-driven delivery pipeline
Model each delivery as a finite state machine (FSM). Every transition writes an immutable, append-only audit record with the state, worker identity, timestamp, and relevant metadata.
RECEIVED ──► ENCRYPTED ──► LEASED ──► ATTEMPTED ──► DELIVERED
▲ │
│ ├── retryable failure ──► RETRIED ──┐
└────────────┼───────────────────────────────────┘
│ (re-queued, then leased again)
└── retries exhausted / non-retryable ──► DLQ
State descriptions
| State | Scope | Trigger |
|---|---|---|
RECEIVED | Ingestion tier | Payload validated and durably persisted. |
ENCRYPTED | Security / KMS layer | Payload encrypted (for example with AES-256-GCM) before it is placed in queue storage. |
LEASED | Queue / worker pool | A worker claims the job via a lease with a fixed time-to-live (TTL). |
ATTEMPTED | Outbound transport | Worker sends the HTTP POST; status code and latency are recorded. |
RETRIED | Scheduler | Non-2xx response or timeout; next attempt time computed with exponential backoff and jitter. |
DELIVERED | Terminal | Receiver returned a 2xx response. |
DLQ | Dead-letter storage | Retry budget exhausted, or a non-retryable error such as most 4xx responses. |
Treat 408 and 429 as retryable rather than as permanent client errors.
Jitter matters. Exponential backoff alone reduces call volume but does not break up synchronized retry clusters, which is why production systems add randomness to the delay.
State transition schema
Store delivery metadata separately from event payloads, so the hot path touches small rows and payload storage stays isolated. The schema below includes the lease columns used later in this guide.
CREATE TABLE webhook_deliveries (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
event_id UUID NOT NULL REFERENCES webhook_events(id),
subscription_id UUID NOT NULL,
current_state VARCHAR(32) NOT NULL,
attempt_count INT NOT NULL DEFAULT 0,
max_attempts INT NOT NULL DEFAULT 5,
next_attempt_at TIMESTAMPTZ,
worker_id VARCHAR(64),
lease_expires_at TIMESTAMPTZ,
fencing_token BIGINT NOT NULL DEFAULT 0,
created_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp()
);
-- Append-only audit trail: one row per state transition
CREATE TABLE webhook_delivery_attempts (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
delivery_id UUID NOT NULL REFERENCES webhook_deliveries(id),
state VARCHAR(32) NOT NULL,
worker_id VARCHAR(64),
fencing_token BIGINT,
http_status_code INT,
latency_ms INT,
error_message TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp()
);
CREATE INDEX idx_deliveries_claimable
ON webhook_deliveries (next_attempt_at)
WHERE current_state IN ('ENCRYPTED', 'RETRIED', 'LEASED');
CREATE INDEX idx_attempts_delivery_time
ON webhook_delivery_attempts (delivery_id, created_at);
On PostgreSQL, gen_random_uuid() is built in from version 13 (earlier versions need the pgcrypto extension). clock_timestamp() returns the actual current time, unlike now(), which is frozen at transaction start. Use it for lease arithmetic. If you run multiple database nodes, prefer the database's clock over worker clocks for all lease comparisons.
To let many workers claim jobs without blocking each other, combine the claim query with FOR UPDATE SKIP LOCKED. Concurrent workers then skip rows another worker has already locked instead of waiting on them.
End-to-end timeline view
With one record per transition, an SRE can open a single delivery and read its history instead of grepping log streams:
[00:00.000] RECEIVED - Payload ingested via gateway node gw-us-east-1a
[00:00.004] ENCRYPTED - Encrypted via KMS (key id: ...8f2a)
[00:00.012] LEASED - Worker wrk-402 claimed lease (30 s), token=1
[00:00.015] ATTEMPTED - POST https://api.client.example/webhook
[00:00.120] FAILED - Endpoint returned HTTP 503 (latency 105 ms)
[00:00.122] RETRIED - Rescheduled for +15 s (exponential backoff + jitter)
[00:15.125] LEASED - Worker wrk-109 claimed lease, token=2
[00:15.128] ATTEMPTED - POST https://api.client.example/webhook
[00:15.210] DELIVERED - HTTP 200 OK (latency 82 ms)
Useful derived metrics from this table include time-in-state percentiles (queue wait is LEASED minus ENCRYPTED), retry rate per subscription, DLQ inflow, and lease-expiry counts. As one webhook-reliability guide notes, a single failed delivery is expected in an at-least-once system, so alert on sustained increases in failure rate, retry rate, queue depth, or dead-letter volume rather than on isolated events.
Build vs. buy
Visual timelines, FSM tracking, and query pipelines are real engineering work. Managed options include Hookdeck, Svix, the open-source and self-hostable Convoy, and relay services such as InstaWebhook. Hookdeck's Event Gateway, for example, provides ingestion, signature verification, deduplication, retries, and replay. Whichever you evaluate, check that it exposes per-attempt history, replay, and DLQ inspection. Also check how it handles duplicate delivery, since that determines how much work stays on your side. hookdeck
2. Preventing Split-Brain Processing: Heartbeats, Leases, and Fencing
When many workers pull from shared queues, you must decide what happens when two workers both believe they own the same job.
Root cause: lease expiry during a pause
Worker A Lease store (database) Worker B
|--- 1. Acquire lease (TTL 10 s, token=1) -->| |
| | |
| [GC pause / VM stall / network partition] | |
| |<-- 2. Lease expired ---|
| |<-- 3. Acquire (token=2)|
| | |--> 5. POST to receiver
|--- 4. Resume; still believes it owns it -->| |
| POST to receiver | |
- Worker A claims a delivery with a 10-second lease.
- Worker A stalls, for example in a long garbage-collection pause or a partition.
- The lease expires without renewal. The system assumes Worker A is dead and makes the job claimable.
- Worker B claims it and sends the HTTP request.
- Worker A wakes up, unaware its lease is gone, and sends the same request.
The receiver now sees two requests for one event. This is the failure mode described in Martin Kleppmann's analysis of distributed locking. He calls leases a flawed basis for correctness unless every write is checked with a fencing token, a number that increases with each lock acquisition, which the storage service uses to reject stale writers. dnitza
Fencing tokens and atomic lease renewal
1. Acquire the lease and increment the token atomically
UPDATE webhook_deliveries
SET current_state = 'LEASED',
worker_id = 'worker-node-102',
lease_expires_at = clock_timestamp() + INTERVAL '10 seconds',
fencing_token = fencing_token + 1,
updated_at = clock_timestamp()
WHERE id = 'd82c4f74-32e7-4b11-9a99-923f6e9112a2'
AND (
current_state IN ('ENCRYPTED', 'RETRIED')
OR (current_state = 'LEASED' AND lease_expires_at < clock_timestamp())
)
RETURNING fencing_token;
If the query returns no row, another worker holds a live lease and this worker must walk away. Because the increment and the claim happen in one statement, tokens are strictly increasing per delivery.
2. Heartbeat loop
While a job runs, a background thread extends the lease. The WHERE clause checks both the worker ID and the token, so a stale worker cannot renew a lease it has lost.
import threading
class WorkerHeartbeat:
def __init__(self, db, delivery_id, worker_id, fencing_token,
interval_sec=3, lease_sec=10):
self.db = db
self.delivery_id = delivery_id
self.worker_id = worker_id
self.fencing_token = fencing_token
self.interval = interval_sec
self.lease_sec = lease_sec
self._stop = threading.Event()
self.lease_lost = threading.Event() # checked by the main work loop
self._thread = threading.Thread(target=self._loop, daemon=True)
def start(self):
self._thread.start()
def stop(self):
self._stop.set()
self._thread.join()
def _loop(self):
# Event.wait returns True when stop is set, ending the loop promptly
while not self._stop.wait(self.interval):
rows = self.db.execute(
"""
UPDATE webhook_deliveries
SET lease_expires_at = clock_timestamp() + make_interval(secs => %s)
WHERE id = %s
AND worker_id = %s
AND fencing_token = %s
AND current_state = 'LEASED'
""",
(self.lease_sec, self.delivery_id, self.worker_id, self.fencing_token),
).rowcount
if rows == 0:
# Lease lost or reassigned. Signal the worker; do not just raise here.
self.lease_lost.set()
break
Two details matter. First, check the affected row count, not just whether the query succeeded, because an UPDATE that matches nothing still succeeds. Second, an exception raised in a background thread does not stop the main worker. The thread must set a flag or cancel the in-flight request, and the main loop must check it before every side effect. Run heartbeats at roughly a third of the lease TTL (3 seconds for a 10-second lease) so one missed beat doesn't expire the lease.
3. Guard every state write with the token
UPDATE webhook_deliveries
SET current_state = 'DELIVERED',
updated_at = clock_timestamp()
WHERE id = %s
AND fencing_token = %s -- fails if another worker has since taken over
AND worker_id = %s;
If this updates zero rows, the worker has been fenced out. It must discard its result and write nothing further. The same check should apply to inserts into the audit table.
What fencing does not protect: the receiver
Fencing tokens work when the resource being written can check them, as your database can. The receiver's HTTP endpoint cannot. If Worker A has already sent its request, or sends it in the instant before it discovers it lost the lease, no token check can recall it. Fencing also does not give mutual exclusion. Kleppmann has agreed that fencing tokens would not protect against two processes both believing they hold the lock, with the lower token reaching the shared resource before the higher one. surfingcomplexity
So the honest guarantee of a lease-and-fencing design is this:
- Your internal state stays consistent. Only the current lease holder can record outcomes.
- The window for duplicate sends is minimized. Check
lease_lostand re-validate the lease immediately before sending. - Duplicate delivery to receivers remains possible. Treat delivery as at-least-once.
The correct complement is idempotency. Send a stable event or delivery ID with every request, for example in a header, identical across retries and across workers. Document that receivers must deduplicate on it. Webhook providers behave this way in practice, which is why guidance for webhook consumers consistently stresses idempotent handlers.
Comparison of concurrency-control strategies
| Optimistic locking | Redis single-node lock | Redlock (multi-node Redis) | Database lease + fencing token | |
|---|---|---|---|---|
| Main risk | High retry rate under heavy contention | Lock lost on failover; no fencing | Relies on timing and clock assumptions | Extra database write load |
| Fencing support | Version column can act as one | None built in | None built in; critics note no obvious way to add monotonic tokens | Yes, if the token is checked on every write |
| Safe for correctness-critical work | Yes, for the guarded row | No | Disputed (see below) | Yes, for state in the same database |
| Complexity | Low | Low | Medium | Medium to high |
| Best for | Low-contention updates | Efficiency-only locks (avoiding duplicate work) | Efficiency-style locks | Job queues where correctness matters |
Redlock is contested. Kleppmann argues it is unsuitable where correctness depends on the lock. Its author, Salvatore Sanfilippo (antirez), responded with a rebuttal, and the two disagree on the timing assumptions. His position is that you can generate an incremental ID for each Redlock acquired if you have a linearizable store. Kleppmann's recommendation for correctness-critical locking is a consensus-backed system such as ZooKeeper, or a database with transactional guarantees plus fencing tokens. For efficiency-only locks, where an occasional duplicate is harmless, a single Redis node is usually enough. antirez
Many teams skip external lock services altogether. They use the database as the serialization point (SELECT … FOR UPDATE SKIP LOCKED or the conditional UPDATE above) or a managed queue with a visibility timeout. Amazon SQS works this way: a received message is hidden for a configurable period and becomes visible again if it isn't deleted in time. It is also at-least-once, so the same duplicate-handling caveat applies. SQS dead-letter queues use a maxReceiveCount threshold to move repeatedly failing messages aside while keeping the original message ID.
Architecture Summary Checklist
- Granular state engine: Every delivery moves through explicit states (
RECEIVED,ENCRYPTED,LEASED,ATTEMPTED,RETRIED,DELIVERED,DLQ), with each transition recorded. - Immutable audit log: Each attempt stores a timestamp, worker ID, fencing token, HTTP status, and latency.
- Backoff with jitter: Retries use capped exponential backoff plus randomness, with a defined retry budget and a clear retryable/non-retryable error policy.
- Atomic leases and fencing tokens: Claims are single conditional statements, and every state write is guarded by the token.
- Heartbeats with real abort handling: Renewal checks affected rows, and the main worker loop stops before any further side effect once the lease is lost.
- Idempotency end to end: A stable delivery ID is sent on every attempt, and receivers are told to deduplicate.
- Recoverable DLQ: Dead-lettered events can be inspected and replayed in bulk, with alerts on sustained DLQ growth rather than single failures.