# Reliable Webhook Delivery (/learn/queues/scenarios/webhook-delivery)



## Architecture [#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.

<Mermaid
  chart="flowchart LR
    APP[&#x22;Application&#x22;]
    WQ{{&#x22;webhooks<br/>queue&#x22;}}
    DW[&#x22;Delivery Worker&#x22;]
    EP[&#x22;External Endpoint&#x22;]
    DLQ{{&#x22;webhooks.dlq&#x22;}}

    APP -- enqueue --> WQ
    WQ -- deliver --> DW
    DW -- &#x22;HTTP POST&#x22; --> EP
    EP -. &#x22;2xx → ack&#x22; .-> DW
    DW -. &#x22;4xx/5xx → requeue (delay)&#x22; .-> WQ
    DW -- &#x22;max retries&#x22; --> DLQ

    class APP,DW client
    class WQ,DLQ queue
    class EP external"
/>

*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 [#webhook-publisher]

When an event occurs, enqueue a webhook delivery request with retry policy.

<Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
  <Tab value="Go">
    ```go title="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
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python title="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}")
    ```
  </Tab>

  <Tab value="Node.js">
    ```typescript title="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}`);
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java title="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);
    }
    ```
  </Tab>

  <Tab value="C#">
    ```csharp title="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
        });
    }
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin title="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
        ))
    }
    ```
  </Tab>

  <Tab value="C++">
    ```cpp title="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);
    }
    ```
  </Tab>

  <Tab value="Rust">
    ```rust title="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(())
    }
    ```
  </Tab>

  <Tab value="Ruby">
    ```ruby title="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
    ```
  </Tab>

  <Tab value="Elixir">
    ```elixir title="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
    ```
  </Tab>
</Tabs>

## Delivery Worker with Exponential Backoff [#delivery-worker-with-exponential-backoff]

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

<Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
  <Tab value="Go">
    ```go title="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()
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python title="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)
    ```
  </Tab>

  <Tab value="Node.js">
    ```typescript title="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));
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java title="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);
    }
    ```
  </Tab>

  <Tab value="C#">
    ```csharp title="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);
    }
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin title="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)
    }
    ```
  </Tab>

  <Tab value="C++">
    ```cpp title="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));
    }
    ```
  </Tab>

  <Tab value="Rust">
    ```rust title="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);
        }
    }
    ```
  </Tab>

  <Tab value="Ruby">
    ```ruby title="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
    ```
  </Tab>

  <Tab value="Elixir">
    ```elixir title="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
    ```
  </Tab>
</Tabs>

## Advanced Configuration [#advanced-configuration]

<Accordions>
  <Accordion title="Delivery Timeout">
    Set a short HTTP timeout (10-30s) per delivery attempt. Long-running endpoints should not block the worker from processing other webhooks.
  </Accordion>

  <Accordion title="Deduplication">
    Use `eventId` as an idempotency key on the receiving end. The same webhook may be delivered more than once if the ack fails after successful delivery.
  </Accordion>

  <Accordion title="Signature Verification">
    Include an HMAC signature in the webhook payload so recipients can verify the sender. Store the signing secret in message tags or metadata.
  </Accordion>
</Accordions>

## Next Steps [#next-steps]

<Cards>
  <Card title="Retry with Backoff" href="/learn/queues/how-to/retry-with-backoff" description="Retry strategies for queue processing." />

  <Card title="Dead Letter Queue" href="/learn/queues/tutorials/dead-letter-queue" description="Automatic DLQ routing." />
</Cards>
