Reliable Webhook Delivery
Guarantee webhook delivery with retry, exponential backoff, and dead letter handling.
Architecture
An application publishes webhook events to a queue. Delivery workers attempt to POST the payload to the target URL. Failed deliveries are retried with exponential backoff using delayed requeue. Permanently failed webhooks land in a Dead Letter Queue (DLQ) for manual investigation.
The worker pulls each webhook, POSTs to the endpoint, acks on success, requeues with an exponential delay on transient failure, and routes to the DLQ once retries are exhausted.
Webhook Publisher
When an event occurs, enqueue a webhook delivery request with retry policy.
type WebhookEvent struct {
EventID string `json:"eventId"`
EventType string `json:"eventType"`
URL string `json:"url"`
Payload string `json:"payload"`
Attempt int `json:"attempt"`
}
func publishWebhook(ctx context.Context, client *kubemq.Client, event WebhookEvent) error {
body, _ := json.Marshal(event)
msg := kubemq.NewQueueMessage().
SetChannel("webhooks").
SetBody(body).
SetMetadata(event.EventType).
SetTags(map[string]string{
"eventId": event.EventID,
"eventType": event.EventType,
"url": event.URL,
}).
SetMaxReceiveCount(5).
SetMaxReceiveQueue("webhooks.dlq").
SetExpirationSeconds(3600)
_, err := client.SendQueueMessage(ctx, msg)
return err
}import json
def publish_webhook(client, event):
result = client.send_queue_message(
QueueMessage(
channel="webhooks",
body=json.dumps(event).encode(),
metadata=event["eventType"],
tags={
"eventId": event["eventId"],
"eventType": event["eventType"],
"url": event["url"],
},
max_receive_count=5,
max_receive_queue="webhooks.dlq",
expiration_in_seconds=3600,
)
)
print(f"Webhook {event['eventId']} enqueued: id={result.id}")async function publishWebhook(client: KubeMQClient, event: WebhookEvent) {
const result = await client.sendQueueMessage(
createQueueMessage({
channel: 'webhooks',
body: JSON.stringify(event),
metadata: event.eventType,
tags: { eventId: event.eventId, eventType: event.eventType, url: event.url },
policy: {
maxReceiveCount: 5,
maxReceiveQueue: 'webhooks.dlq',
expirationSeconds: 3600,
},
}),
);
console.log(`Webhook ${event.eventId} enqueued: id=${result.messageId}`);
}public void publishWebhook(QueuesClient client, WebhookEvent event) throws Exception {
QueueMessage msg = QueueMessage.builder()
.channel("webhooks")
.body(objectMapper.writeValueAsBytes(event))
.metadata(event.getEventType())
.tags(Map.of(
"eventId", event.getEventId(),
"eventType", event.getEventType(),
"url", event.getUrl()))
.maxReceiveCount(5)
.maxReceiveQueue("webhooks.dlq")
.expirationSeconds(3600)
.build();
client.sendQueueMessage(msg);
}async Task PublishWebhook(QueuesClient client, WebhookEvent evt)
{
await client.SendQueueMessageAsync(new QueueMessage
{
Channel = "webhooks",
Body = JsonSerializer.SerializeToUtf8Bytes(evt),
Metadata = evt.EventType,
Tags = new Dictionary<string, string>
{
["eventId"] = evt.EventId,
["eventType"] = evt.EventType,
["url"] = evt.Url
},
MaxReceiveCount = 5,
MaxReceiveQueue = "webhooks.dlq",
ExpirationSeconds = 3600
});
}suspend fun publishWebhook(client: QueuesClient, event: WebhookEvent) {
client.sendQueueMessage(QueueMessage(
channel = "webhooks",
body = Json.encodeToString(event).toByteArray(),
metadata = event.eventType,
tags = mapOf(
"eventId" to event.eventId,
"eventType" to event.eventType,
"url" to event.url),
maxReceiveCount = 5,
maxReceiveQueue = "webhooks.dlq",
expirationSeconds = 3600
))
}void publishWebhook(kubemq::QueuesClient& client, const WebhookEvent& event) {
kubemq::QueueMessage msg;
msg.channel = "webhooks";
msg.body = event.toJson();
msg.metadata = event.eventType;
msg.tags = {
{"eventId", event.eventId},
{"eventType", event.eventType},
{"url", event.url}};
msg.maxReceiveCount = 5;
msg.maxReceiveQueue = "webhooks.dlq";
msg.expirationSeconds = 3600;
client.sendQueueMessage(msg);
}use kubemq::prelude::*;
use kubemq::QueueMessageBuilder;
use std::collections::HashMap;
async fn publish_webhook(
client: &KubemqClient,
event: &WebhookEvent,
) -> kubemq::Result<()> {
let mut tags = HashMap::new();
tags.insert("eventId".to_string(), event.event_id.clone());
tags.insert("eventType".to_string(), event.event_type.clone());
tags.insert("url".to_string(), event.url.clone());
let msg = QueueMessageBuilder::new()
.channel("webhooks")
.body(serde_json::to_vec(event).unwrap())
.metadata(&event.event_type)
.tags(tags)
.max_receive_count(5)
.max_receive_queue("webhooks.dlq")
.expiration_seconds(3600)
.build();
let result = client.send_queue_message(msg).await?;
println!("Webhook {} enqueued: id={}", event.event_id, result.message_id);
Ok(())
}require 'kubemq'
require 'json'
def publish_webhook(client, event)
policy = KubeMQ::Queues::QueueMessagePolicy.new(
max_receive_count: 5,
max_receive_queue: 'webhooks.dlq',
expiration_seconds: 3600
)
msg = KubeMQ::Queues::QueueMessage.new(
channel: 'webhooks',
metadata: event['eventType'],
body: event.to_json,
tags: {
'eventId' => event['eventId'],
'eventType' => event['eventType'],
'url' => event['url']
},
policy: policy
)
result = client.send_queue_message(msg)
puts "Webhook #{event['eventId']} enqueued: id=#{result.id}"
enddef publish_webhook(client, event) do
msg =
KubeMQ.QueueMessage.new(
channel: "webhooks",
metadata: event.event_type,
body: Jason.encode!(event),
tags: %{
"eventId" => event.event_id,
"eventType" => event.event_type,
"url" => event.url
},
policy:
KubeMQ.QueuePolicy.new(
max_receive_count: 5,
max_receive_queue: "webhooks.dlq",
expiration_seconds: 3600
)
)
{:ok, result} = KubeMQ.Client.send_queue_message(client, msg)
IO.puts("Webhook #{event.event_id} enqueued: id=#{result.message_id}")
endDelivery Worker with Exponential Backoff
Attempt HTTP delivery. On failure, requeue with exponential delay.
for {
resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
Channel: "webhooks",
MaxItems: 10,
WaitTimeoutSeconds: 10,
VisibilitySeconds: 30,
})
if err != nil {
log.Printf("Poll error: %v", err)
time.Sleep(5 * time.Second)
continue
}
for _, m := range resp.Messages {
var event WebhookEvent
json.Unmarshal(m.Message.Body, &event)
httpResp, err := http.Post(event.URL, "application/json",
bytes.NewReader([]byte(event.Payload)))
if err == nil && httpResp.StatusCode >= 200 && httpResp.StatusCode < 300 {
log.Printf("Webhook %s delivered to %s", event.EventID, event.URL)
continue
}
attempt := event.Attempt + 1
if attempt >= 5 {
log.Printf("Webhook %s failed permanently, sending to DLQ", event.EventID)
continue
}
delay := 1 << attempt // 2, 4, 8, 16 seconds
event.Attempt = attempt
retryBody, _ := json.Marshal(event)
client.SendQueueMessage(ctx, kubemq.NewQueueMessage().
SetChannel("webhooks").
SetBody(retryBody).
SetDelaySeconds(delay).
SetMaxReceiveCount(5).
SetMaxReceiveQueue("webhooks.dlq"))
log.Printf("Webhook %s retry %d in %ds", event.EventID, attempt, delay)
}
resp.AckAll()
}import json
import requests
import time
while True:
response = client.receive_queue_messages(
channel="webhooks",
max_messages=10,
wait_timeout_in_seconds=10,
visibility_seconds=30,
)
for msg in response.messages:
event = json.loads(msg.body.decode("utf-8"))
try:
resp = requests.post(event["url"], json=json.loads(event["payload"]), timeout=10)
resp.raise_for_status()
print(f"Webhook {event['eventId']} delivered to {event['url']}")
msg.ack()
except Exception as e:
attempt = event.get("attempt", 0) + 1
if attempt >= 5:
print(f"Webhook {event['eventId']} failed permanently")
msg.requeue("webhooks.dlq")
continue
delay = 2 ** attempt
event["attempt"] = attempt
client.send_queue_message(QueueMessage(
channel="webhooks",
body=json.dumps(event).encode(),
delay_in_seconds=delay,
max_receive_count=5,
max_receive_queue="webhooks.dlq",
))
msg.ack()
print(f"Webhook {event['eventId']} retry {attempt} in {delay}s")
time.sleep(1)while (true) {
const messages = await client.receiveQueueMessages({
channel: 'webhooks',
maxMessages: 10,
waitTimeoutSeconds: 10,
visibilitySeconds: 30,
});
for (const msg of messages) {
const event = JSON.parse(new TextDecoder().decode(msg.body));
try {
const resp = await fetch(event.url, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: event.payload,
});
if (resp.ok) {
console.log(`Webhook ${event.eventId} delivered`);
await msg.ack();
continue;
}
throw new Error(`HTTP ${resp.status}`);
} catch (err) {
const attempt = (event.attempt || 0) + 1;
if (attempt >= 5) {
console.log(`Webhook ${event.eventId} failed permanently`);
await msg.requeue('webhooks.dlq');
continue;
}
const delay = Math.pow(2, attempt);
event.attempt = attempt;
await client.sendQueueMessage(
createQueueMessage({
channel: 'webhooks',
body: JSON.stringify(event),
policy: { delaySeconds: delay, maxReceiveCount: 5, maxReceiveQueue: 'webhooks.dlq' },
}),
);
await msg.ack();
console.log(`Webhook ${event.eventId} retry ${attempt} in ${delay}s`);
}
}
await new Promise((r) => setTimeout(r, 1000));
}while (true) {
ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
ReceiveQueueMessagesRequest.builder()
.channel("webhooks").maxMessages(10)
.waitTimeoutSeconds(10).visibilitySeconds(30).build());
for (QueueMessageReceived msg : response.getMessages()) {
WebhookEvent event = objectMapper.readValue(msg.getBody(), WebhookEvent.class);
try {
HttpResponse<String> resp = httpClient.send(
HttpRequest.newBuilder().uri(URI.create(event.getUrl()))
.POST(HttpRequest.BodyPublishers.ofString(event.getPayload())).build(),
HttpResponse.BodyHandlers.ofString());
if (resp.statusCode() >= 200 && resp.statusCode() < 300) {
msg.ack();
continue;
}
throw new Exception("HTTP " + resp.statusCode());
} catch (Exception e) {
int attempt = event.getAttempt() + 1;
if (attempt >= 5) { msg.requeue("webhooks.dlq"); continue; }
int delay = (int) Math.pow(2, attempt);
event.setAttempt(attempt);
client.sendQueueMessage(QueueMessage.builder()
.channel("webhooks").body(objectMapper.writeValueAsBytes(event))
.delaySeconds(delay).build());
msg.ack();
}
}
Thread.sleep(1000);
}while (true)
{
var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
{
Channel = "webhooks", MaxMessages = 10,
WaitTimeoutSeconds = 10, VisibilitySeconds = 30,
});
foreach (var msg in response.Messages)
{
var evt = JsonSerializer.Deserialize<WebhookEvent>(msg.Body.Span);
try
{
var resp = await httpClient.PostAsync(evt.Url,
new StringContent(evt.Payload, Encoding.UTF8, "application/json"));
resp.EnsureSuccessStatusCode();
await msg.AckAsync();
}
catch
{
var attempt = evt.Attempt + 1;
if (attempt >= 5) { await msg.ReQueueAsync("webhooks.dlq"); continue; }
evt.Attempt = attempt;
var delay = (int)Math.Pow(2, attempt);
await client.SendQueueMessageAsync(new QueueMessage
{
Channel = "webhooks",
Body = JsonSerializer.SerializeToUtf8Bytes(evt),
DelaySeconds = delay
});
await msg.AckAsync();
}
}
await Task.Delay(1000);
}while (true) {
val response = client.receiveQueueMessages(
channel = "webhooks", maxMessages = 10,
waitTimeoutSeconds = 10, visibilitySeconds = 30)
for (msg in response.messages) {
val event = Json.decodeFromString<WebhookEvent>(String(msg.body))
try {
val resp = httpClient.post(event.url) { setBody(event.payload) }
if (resp.status.value in 200..299) { msg.ack(); continue }
throw Exception("HTTP ${resp.status}")
} catch (e: Exception) {
val attempt = event.attempt + 1
if (attempt >= 5) { msg.requeue("webhooks.dlq"); continue }
val delay = 2.0.pow(attempt.toDouble()).toInt()
val retryEvent = event.copy(attempt = attempt)
client.sendQueueMessage(QueueMessage(
channel = "webhooks",
body = Json.encodeToString(retryEvent).toByteArray(),
delaySeconds = delay))
msg.ack()
}
}
delay(1000)
}while (true) {
auto response = client.receiveQueueMessages("webhooks", 10, 10, false, 30);
for (const auto& msg : response.messages) {
auto event = WebhookEvent::fromJson(msg.body);
auto httpResp = httpPost(event.url, event.payload);
if (httpResp.statusCode >= 200 && httpResp.statusCode < 300) {
msg.ack();
continue;
}
int attempt = event.attempt + 1;
if (attempt >= 5) { msg.requeue("webhooks.dlq"); continue; }
int delaySec = static_cast<int>(std::pow(2, attempt));
event.attempt = attempt;
kubemq::QueueMessage retryMsg;
retryMsg.channel = "webhooks";
retryMsg.body = event.toJson();
retryMsg.delaySeconds = delaySec;
client.sendQueueMessage(retryMsg);
msg.ack();
}
std::this_thread::sleep_for(std::chrono::seconds(1));
}use kubemq::prelude::*;
use kubemq::{PollRequest, QueueMessageBuilder};
let mut receiver = client.new_queue_downstream_receiver().await?;
loop {
let poll = PollRequest {
channel: "webhooks".to_string(),
max_items: 10,
wait_timeout_seconds: 10,
auto_ack: false,
};
let response = receiver.poll(poll).await?;
for m in &response.messages {
let event: WebhookEvent = serde_json::from_slice(&m.message.body).unwrap();
let resp = reqwest::Client::new()
.post(&event.url)
.body(event.payload.clone())
.send()
.await;
if let Ok(r) = resp {
if r.status().is_success() {
m.ack().await?;
println!("Webhook {} delivered to {}", event.event_id, event.url);
continue;
}
}
let attempt = event.attempt + 1;
if attempt >= 5 {
// Route to the dead-letter queue after exhausting retries.
m.re_queue("webhooks.dlq").await?;
println!("Webhook {} failed permanently, sent to DLQ", event.event_id);
continue;
}
// Re-send with an exponential delay, then ack the original.
let delay = 1 << attempt; // 2, 4, 8, 16 seconds
let mut retry = event.clone();
retry.attempt = attempt;
let retry_msg = QueueMessageBuilder::new()
.channel("webhooks")
.body(serde_json::to_vec(&retry).unwrap())
.delay_seconds(delay)
.max_receive_count(5)
.max_receive_queue("webhooks.dlq")
.build();
client.send_queue_message(retry_msg).await?;
m.ack().await?;
println!("Webhook {} retry {} in {}s", event.event_id, attempt, delay);
}
}require 'kubemq'
require 'net/http'
require 'json'
receiver = client.create_downstream_receiver
loop do
request = KubeMQ::Queues::QueuePollRequest.new(
channel: 'webhooks', max_items: 10, wait_timeout: 10
)
response = receiver.poll(request)
next if response.error?
response.messages.each do |m|
event = JSON.parse(m.body)
begin
uri = URI(event['url'])
http_resp = Net::HTTP.post(uri, event['payload'], 'Content-Type' => 'application/json')
raise "HTTP #{http_resp.code}" unless http_resp.is_a?(Net::HTTPSuccess)
m.ack
puts "Webhook #{event['eventId']} delivered to #{event['url']}"
rescue StandardError
attempt = (event['attempt'] || 0) + 1
if attempt >= 5
# Route to the dead-letter queue after exhausting retries:
# re-send the event to the DLQ channel, then ack the original.
client.send_queue_message(KubeMQ::Queues::QueueMessage.new(
channel: 'webhooks.dlq', body: event.to_json
))
m.ack
puts "Webhook #{event['eventId']} failed permanently, sent to DLQ"
next
end
# Re-send with an exponential delay, then ack the original.
delay = 2**attempt
event['attempt'] = attempt
retry_policy = KubeMQ::Queues::QueueMessagePolicy.new(
delay_seconds: delay, max_receive_count: 5, max_receive_queue: 'webhooks.dlq'
)
client.send_queue_message(KubeMQ::Queues::QueueMessage.new(
channel: 'webhooks', body: event.to_json, policy: retry_policy
))
m.ack
puts "Webhook #{event['eventId']} retry #{attempt} in #{delay}s"
end
end
enddefp worker_loop(client) do
case KubeMQ.Client.poll_queue(client,
channel: "webhooks",
max_items: 10,
wait_timeout: 10_000
) do
{:ok, poll} ->
Enum.each(poll.messages, fn m ->
event = Jason.decode!(m.body)
attempt = Map.get(event, "attempt", 0) + 1
seq = m.attributes.sequence
case deliver(event["url"], event["payload"]) do
:ok ->
KubeMQ.PollResponse.ack_range(poll, [seq])
IO.puts("Webhook #{event["eventId"]} delivered to #{event["url"]}")
:error when attempt >= 5 ->
# Route to the dead-letter queue after exhausting retries.
KubeMQ.PollResponse.requeue_range(poll, [seq], "webhooks.dlq")
IO.puts("Webhook #{event["eventId"]} failed permanently, sent to DLQ")
:error ->
# Re-send with an exponential delay, then ack the original.
delay = trunc(:math.pow(2, attempt))
retry = Map.put(event, "attempt", attempt)
KubeMQ.Client.send_queue_message(client,
KubeMQ.QueueMessage.new(
channel: "webhooks",
body: Jason.encode!(retry),
policy:
KubeMQ.QueuePolicy.new(
delay_seconds: delay,
max_receive_count: 5,
max_receive_queue: "webhooks.dlq"
)
)
)
KubeMQ.PollResponse.ack_range(poll, [seq])
IO.puts("Webhook #{event["eventId"]} retry #{attempt} in #{delay}s")
end
end)
{:error, err} ->
IO.puts("Poll error: #{err.message}")
end
worker_loop(client)
endAdvanced Configuration
Next Steps
Was this page helpful?