KubeMQ
LearnQueuesScenarios

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.

webhook_publisher.go
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
}
webhook_publisher.py
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}")
webhook_publisher.ts
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}`);
}
WebhookPublisher.java
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);
}
WebhookPublisher.cs
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
    });
}
WebhookPublisher.kt
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
    ))
}
webhook_publisher.cpp
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);
}
webhook_publisher.rs
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(())
}
webhook_publisher.rb
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}"
end
webhook_publisher.exs
def 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}")
end

Delivery Worker with Exponential Backoff

Attempt HTTP delivery. On failure, requeue with exponential delay.

delivery_worker.go
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()
}
delivery_worker.py
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)
delivery_worker.ts
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));
}
DeliveryWorker.java
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);
}
DeliveryWorker.cs
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);
}
DeliveryWorker.kt
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)
}
delivery_worker.cpp
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));
}
delivery_worker.rs
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);
    }
}
delivery_worker.rb
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
end
delivery_worker.exs
defp 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)
end

Advanced Configuration

Next Steps

Was this page helpful?

On this page