# Queues (/connectors/cloudevents/how-to/queues)



The CloudEvents connector exposes KubeMQ's durable **queue** pattern over plain HTTP: producers `POST` a CloudEvent to a queue, and one consumer at a time pulls messages off it with at-least-once delivery.

## Overview [#overview]

A queue is a **durable, FIFO work channel**. Unlike events (fire-and-forget pub/sub), a queue message is stored until a consumer receives it, and is delivered to exactly one consumer in the group — making queues the right fit for distributing work across a pool of workers.

The connector maps three control operations onto HTTP:

| Operation | Endpoint                 | Body                                      | Purpose                               |
| --------- | ------------------------ | ----------------------------------------- | ------------------------------------- |
| Send      | `POST /ce/queue/send`    | CloudEvent (structured or binary)         | Enqueue one message                   |
| Receive   | `POST /ce/queue/receive` | *(ignored)* — all params via query string | Pull (and consume) messages           |
| Ack all   | `POST /ce/queue/ack_all` | *(ignored)* — all params via query string | Drain all pending messages atomically |

`receive` and `ack_all` are **control operations**: the request body is ignored and every parameter is supplied as a query string. To inspect messages without consuming them, set `is_peek=true` on a `receive` call (peek).

## How it works [#how-it-works]

A producer enqueues CloudEvents; a worker polls the queue and the connector hands each stored message to exactly one consumer.

<Mermaid
  chart="`
graph LR
PROD[&#x22;Producer&#x22;]
CE[&#x22;CloudEvents<br/>connector&#x22;]
Q{{&#x22;Queue channel<br/>work-queue&#x22;}}
W1[&#x22;Worker A&#x22;]
W2[&#x22;Worker B&#x22;]

PROD -- &#x22;POST /ce/queue/send&#x22; --> CE
CE --> Q
W1 -- &#x22;POST /ce/queue/receive&#x22; --> CE
W2 -- &#x22;POST /ce/queue/receive&#x22; --> CE
Q -- &#x22;one message<br/>per consumer&#x22; --> W1
Q -. &#x22;next message&#x22; .-> W2

class PROD,W1,W2 client
class CE connector
class Q queue
`"
/>

*Messages are durably stored on the queue channel and delivered to one worker at a time.*

## Send and receive [#send-and-receive]

Send a CloudEvent to a queue with `POST /ce/queue/send` (returns HTTP 202), then pull it back with `POST /ce/queue/receive`. The receive call is a control operation — pass `channel`, `client_id`, `max_messages`, and `wait_timeout` as query parameters.

<Tabs groupId="language" items="['curl','C#','Go','Java','JavaScript','Python','Ruby','Rust']">
  <Tab value="curl">
    ```bash
    # Send a message to the queue (structured CloudEvent)
    curl -X POST http://localhost:9090/ce/queue/send \
      -H "Content-Type: application/cloudevents+json" \
      -d '{
        "specversion": "1.0",
        "type": "com.example.job.submit",
        "source": "job-scheduler",
        "subject": "work-queue",
        "data": {"job_type": "report", "params": {"month": "2026-03"}}
      }'

    # Receive (consume) up to 10 messages, waiting up to 10s
    curl -X POST "http://localhost:9090/ce/queue/receive?channel=work-queue&client_id=worker-1&max_messages=10&wait_timeout=10"
    ```
  </Tab>

  <Tab value="C#">
    ```csharp
    using CloudNative.CloudEvents;
    using CloudNative.CloudEvents.SystemTextJson;
    using System.Net.Http.Headers;
    using System.Text.Json;

    static string ServerUrl() => Environment.GetEnvironmentVariable("KUBEMQ_CE_URL") ?? "http://localhost:9090";

    var base_ = ServerUrl();
    var channel = "csharp-ce-queues.basic";
    var clientId = "kubemq-ce-csharp-worker";
    var formatter = new JsonEventFormatter();
    using var httpClient = new HttpClient();

    Console.WriteLine($"Sending 3 messages to queue '{channel}':");
    for (int i = 1; i <= 3; i++)
    {
        var ev = new CloudEvent {
            Id = Guid.NewGuid().ToString(),
            Type = "com.kubemq.examples.queues.task",
            Source = new Uri("urn:kubemq-ce-csharp-example"),
            Subject = channel,
            DataContentType = "application/json",
            Data = new { task_id = i, task = "process-item" },
        };
        var bytes = formatter.EncodeStructuredModeMessage(ev, out var ct);
        using var c = new ByteArrayContent(bytes.ToArray());
        c.Headers.ContentType = MediaTypeHeaderValue.Parse(ct.ToString());
        var r = await httpClient.PostAsync($"{base_}/ce/queue/send", c);
        var j = JsonSerializer.Deserialize<JsonElement>(await r.Content.ReadAsStringAsync());
        Console.WriteLine($"  Sent task {i}: is_error={j.GetProperty("is_error")}");
    }

    Console.WriteLine("\nReceiving 3 messages:");
    for (int i = 0; i < 3; i++)
    {
        var url = $"{base_}/ce/queue/receive?channel={Uri.EscapeDataString(channel)}&client_id={clientId}&max_messages=1&wait_timeout=5";
        var r = await httpClient.PostAsync(url, null);
        var j = JsonSerializer.Deserialize<JsonElement>(await r.Content.ReadAsStringAsync());
        if (j.GetProperty("is_error").GetBoolean()) {
            Console.WriteLine($"  Error: {j.GetProperty("message")}"); continue;
        }
        if (j.TryGetProperty("data", out var data) && data.TryGetProperty("messages", out var msgs))
        {
            foreach (var m in msgs.EnumerateArray())
                Console.WriteLine($"  Received: type={m.GetProperty("type")} data={m.GetProperty("data")}");
        }
    }
    ```
  </Tab>

  <Tab value="Go">
    ```go
    package main

    import (
    	"encoding/json"
    	"fmt"
    	"log"
    	"net/http"
    	"os"
    	"strings"

    	cloudevents "github.com/cloudevents/sdk-go/v2"
    )

    func serverURL() string {
    	if u := os.Getenv("KUBEMQ_CE_URL"); u != "" {
    		return u
    	}
    	return "http://localhost:9090"
    }

    type CEResponse struct {
    	IsError bool            `json:"is_error"`
    	Message string          `json:"message"`
    	Data    json.RawMessage `json:"data"`
    }

    type QueueReceiveData struct {
    	MessagesReceived int                      `json:"messages_received"`
    	Messages         []map[string]interface{} `json:"messages"`
    }

    func sendToQueue(base, channel string, i int) {
    	event := cloudevents.NewEvent()
    	event.SetType("com.kubemq.examples.queues.task")
    	event.SetSource("kubemq-ce-go-example")
    	event.SetSubject(channel)
    	_ = event.SetData(cloudevents.ApplicationJSON, map[string]interface{}{
    		"task_id": i, "task": "process-item", "priority": "normal",
    	})

    	body, _ := json.Marshal(event)
    	req, _ := http.NewRequest("POST", base+"/ce/queue/send", strings.NewReader(string(body)))
    	req.Header.Set("Content-Type", "application/cloudevents+json")
    	resp, err := http.DefaultClient.Do(req)
    	if err != nil {
    		log.Fatal("queue send:", err)
    	}
    	defer resp.Body.Close()
    	var result CEResponse
    	_ = json.NewDecoder(resp.Body).Decode(&result)
    	fmt.Printf("Sent task %d: status=%d is_error=%v\n", i, resp.StatusCode, result.IsError)
    }

    func receiveFromQueue(base, channel, clientID string) {
    	url := fmt.Sprintf("%s/ce/queue/receive?channel=%s&client_id=%s&max_messages=1&wait_timeout=5",
    		base, channel, clientID)
    	req, _ := http.NewRequest("POST", url, nil)
    	resp, err := http.DefaultClient.Do(req)
    	if err != nil {
    		log.Fatal("queue receive:", err)
    	}
    	defer resp.Body.Close()
    	var result CEResponse
    	_ = json.NewDecoder(resp.Body).Decode(&result)
    	var data QueueReceiveData
    	_ = json.Unmarshal(result.Data, &data)
    	for _, msg := range data.Messages {
    		fmt.Printf("Received: type=%v data=%v\n", msg["type"], msg["data"])
    	}
    }

    func main() {
    	base := serverURL()
    	channel := "go-ce-queues.basic"
    	clientID := "kubemq-ce-go-worker"
    	for i := 1; i <= 3; i++ {
    		sendToQueue(base, channel, i)
    	}
    	for i := 0; i < 3; i++ {
    		receiveFromQueue(base, channel, clientID)
    	}
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java
    import com.fasterxml.jackson.databind.ObjectMapper;
    import io.cloudevents.CloudEvent;
    import io.cloudevents.core.builder.CloudEventBuilder;
    import io.cloudevents.core.format.EventFormat;
    import io.cloudevents.core.provider.EventFormatProvider;
    import io.cloudevents.jackson.JsonFormat;

    import java.net.URI;
    import java.net.http.HttpClient;
    import java.net.http.HttpRequest;
    import java.net.http.HttpResponse;
    import java.time.OffsetDateTime;
    import java.util.List;
    import java.util.Map;
    import java.util.UUID;

    public class Main {
        static String serverUrl() {
            String u = System.getenv("KUBEMQ_CE_URL");
            return (u != null && !u.isEmpty()) ? u : "http://localhost:9090";
        }
        static final ObjectMapper MAPPER = new ObjectMapper();

        public static void main(String[] args) throws Exception {
            String base = serverUrl();
            String channel = "java-ce-queues.basic";
            String clientId = "kubemq-ce-java-worker";

            EventFormatProvider.getInstance().registerFormat(new JsonFormat());
            EventFormat format = EventFormatProvider.getInstance().resolveFormat(JsonFormat.CONTENT_TYPE);
            HttpClient httpClient = HttpClient.newHttpClient();

            for (int i = 1; i <= 3; i++) {
                CloudEvent event = CloudEventBuilder.v1()
                        .withId(UUID.randomUUID().toString())
                        .withType("com.kubemq.examples.queues.task")
                        .withSource(URI.create("kubemq-ce-java-example"))
                        .withSubject(channel)
                        .withDataContentType("application/json")
                        .withTime(OffsetDateTime.now())
                        .withData("application/json",
                                MAPPER.writeValueAsBytes(Map.of("task_id", i, "task", "process-item")))
                        .build();
                httpClient.send(
                        HttpRequest.newBuilder().uri(URI.create(base + "/ce/queue/send"))
                                .POST(HttpRequest.BodyPublishers.ofByteArray(format.serialize(event)))
                                .header("Content-Type", "application/cloudevents+json").build(),
                        HttpResponse.BodyHandlers.ofString());
            }

            for (int i = 0; i < 3; i++) {
                String url = base + "/ce/queue/receive?channel=" + channel
                        + "&client_id=" + clientId + "&max_messages=1&wait_timeout=5";
                HttpResponse<String> resp = httpClient.send(
                        HttpRequest.newBuilder().uri(URI.create(url))
                                .POST(HttpRequest.BodyPublishers.noBody()).build(),
                        HttpResponse.BodyHandlers.ofString());
                Map<?, ?> r = MAPPER.readValue(resp.body(), Map.class);
                Map<?, ?> data = (Map<?, ?>) r.get("data");
                List<?> msgs = data != null ? (List<?>) data.get("messages") : List.of();
                for (Object m : msgs) {
                    Map<?, ?> msg = (Map<?, ?>) m;
                    System.out.println("  Received: type=" + msg.get("type") + " data=" + msg.get("data"));
                }
            }
        }
    }
    ```
  </Tab>

  <Tab value="JavaScript">
    ```javascript
    import { CloudEvent, HTTP } from 'cloudevents';

    function serverUrl() {
      return process.env.KUBEMQ_CE_URL ?? 'http://localhost:9090';
    }

    async function main() {
      const base = serverUrl();
      const channel = 'js-ce-queues.basic';
      const clientId = 'kubemq-ce-js-worker';
      const numMessages = 3;

      for (let i = 1; i <= numMessages; i++) {
        const event = new CloudEvent({
          type: 'com.kubemq.examples.queues.task',
          source: 'kubemq-ce-js-example',
          subject: channel,
          datacontenttype: 'application/json',
          data: { task_id: i, task: 'process-item' },
        });
        const msg = HTTP.structured(event);
        const resp = await fetch(`${base}/ce/queue/send`, {
          method: 'POST',
          headers: msg.headers,
          body: msg.body,
        });
        const result = await resp.json();
        console.log(`  Sent task ${i}: status=${resp.status} is_error=${result.is_error}`);
      }

      for (let i = 0; i < numMessages; i++) {
        const url = `${base}/ce/queue/receive?channel=${encodeURIComponent(channel)}&client_id=${clientId}&max_messages=1&wait_timeout=5`;
        const resp = await fetch(url, { method: 'POST' });
        const result = await resp.json();
        for (const m of result.data?.messages ?? []) {
          console.log(`  Received: type=${m.type} data=${JSON.stringify(m.data)}`);
        }
      }
    }

    main().catch(console.error);
    ```
  </Tab>

  <Tab value="Python">
    ```python
    import os
    import requests
    from cloudevents.v1.conversion import to_structured
    from cloudevents.v1.http import CloudEvent


    def server_url() -> str:
        return os.environ.get("KUBEMQ_CE_URL", "http://localhost:9090")


    def main() -> None:
        base = server_url()
        channel = "python-ce-queues.basic"
        client_id = "kubemq-ce-python-worker"
        num_messages = 3

        for i in range(1, num_messages + 1):
            event = CloudEvent(
                attributes={
                    "type": "com.kubemq.examples.queues.task",
                    "source": "kubemq-ce-python-example",
                    "subject": channel,
                    "datacontenttype": "application/json",
                },
                data={"task_id": i, "task": "process-item"},
            )
            headers, body = to_structured(event)
            resp = requests.post(f"{base}/ce/queue/send", data=body,
                                 headers=dict(headers), timeout=10)
            result = resp.json()
            print(f"  Sent task {i}: status={resp.status_code} is_error={result.get('is_error')}")

        for _ in range(num_messages):
            resp = requests.post(
                f"{base}/ce/queue/receive",
                params={
                    "channel": channel,
                    "client_id": client_id,
                    "max_messages": 1,
                    "wait_timeout": 5,
                },
                timeout=10,
            )
            result = resp.json()
            for msg in result.get("data", {}).get("messages", []):
                print(f"  Received: type={msg.get('type')} data={msg.get('data')}")


    if __name__ == "__main__":
        main()
    ```
  </Tab>

  <Tab value="Ruby">
    ```ruby
    require "net/http"
    require "uri"
    require "json"
    require "securerandom"
    require "cloud_events"

    def server_url = ENV.fetch("KUBEMQ_CE_URL", "http://localhost:9090")

    base      = server_url
    channel   = "ruby-ce-queues.basic"
    client_id = "kubemq-ce-ruby-worker"
    sdk       = CloudEvents::HttpBinding.default

    (1..3).each do |i|
      event = CloudEvents::Event::V1.new(
        id: SecureRandom.uuid, type: "com.kubemq.examples.queues.task",
        source: URI("urn:kubemq-ce-ruby-example"), subject: channel,
        spec_version: "1.0",
        data_content_type: CloudEvents::ContentType.new("application/json"),
        data: JSON.generate({ task_id: i, task: "process-item" })
      )
      enc_headers, enc_body = sdk.encode_event(event, structured_format: "json")
      uri = URI("#{base}/ce/queue/send")
      Net::HTTP.start(uri.host, uri.port) do |http|
        req = Net::HTTP::Post.new(uri)
        enc_headers.each { |k, v| req[k] = v }
        req.body = enc_body
        res = http.request(req)
        r = JSON.parse(res.body)
        puts "  Sent task #{i}: status=#{res.code} is_error=#{r['is_error']}"
      end
    end

    3.times do
      uri = URI("#{base}/ce/queue/receive")
      uri.query = URI.encode_www_form(channel: channel, client_id: client_id,
                                       max_messages: 1, wait_timeout: 5)
      Net::HTTP.start(uri.host, uri.port) do |http|
        res = http.request(Net::HTTP::Post.new(uri))
        r = JSON.parse(res.body)
        msgs = r.dig('data', 'messages') || []
        msgs.each { |m| puts "  Received: type=#{m['type']} data=#{m['data']}" }
      end
    end
    ```
  </Tab>

  <Tab value="Rust">
    ```rust
    use cloudevents::{EventBuilder, EventBuilderV10};
    use reqwest::Client;
    use serde_json::{json, Value};
    use std::env;
    use uuid::Uuid;

    fn server_url() -> String {
        env::var("KUBEMQ_CE_URL").unwrap_or_else(|_| "http://localhost:9090".to_string())
    }

    #[tokio::main]
    async fn main() -> Result<(), Box<dyn std::error::Error>> {
        let base = server_url();
        let client = Client::new();
        let channel = "rust-ce-queues.basic";
        let client_id = "kubemq-ce-rust-worker";

        for i in 1..=3u32 {
            let event = EventBuilderV10::new()
                .id(Uuid::new_v4().to_string())
                .ty("com.kubemq.examples.queues.task")
                .source("urn:kubemq-ce-rust-example")
                .subject(channel)
                .data("application/json", json!({"task_id": i, "task": "process-item"}))
                .build()?;
            let body = serde_json::to_string(&event)?;
            let resp = client.post(format!("{}/ce/queue/send", base))
                .header("Content-Type", "application/cloudevents+json")
                .body(body).send().await?;
            let r: Value = resp.json().await?;
            println!("  Sent task {}: is_error={}", i, r["is_error"]);
        }

        for _ in 0..3 {
            let url = format!(
                "{}/ce/queue/receive?channel={}&client_id={}&max_messages=1&wait_timeout=5",
                base, channel, client_id
            );
            let resp = client.post(&url).send().await?;
            let r: Value = resp.json().await?;
            let msgs = r["data"]["messages"].as_array().cloned().unwrap_or_default();
            for m in &msgs {
                println!("  Received: type={} data={}", m["type"], m["data"]);
            }
        }
        Ok(())
    }
    ```
  </Tab>
</Tabs>

### Receive parameters [#receive-parameters]

| Parameter      | Type   | Default      | Description                                                        |
| -------------- | ------ | ------------ | ------------------------------------------------------------------ |
| `channel`      | string | *(required)* | Queue channel name                                                 |
| `client_id`    | string | *(required)* | Client identifier (overridden by auth claims when auth is enabled) |
| `max_messages` | int    | `1`          | Maximum number of messages to receive (1–1000)                     |
| `wait_timeout` | int    | `5`          | Wait timeout in seconds                                            |
| `is_peek`      | bool   | `false`      | If true, peek at messages without consuming them                   |

A successful receive returns HTTP 200. Each message that carries CE tags is reconstructed as a CloudEvent JSON object; messages without CE tags are returned in plain KubeMQ format:

```json
{
  "is_error": false,
  "message": "OK",
  "data": {
    "messages_received": 1,
    "messages": [
      {
        "specversion": "1.0",
        "type": "com.example.job.submit",
        "source": "job-scheduler",
        "subject": "work-queue",
        "id": "550e8400-e29b-41d4-a716-446655440000",
        "time": "2026-03-29T10:30:00Z",
        "data": {"job_type": "report", "params": {"month": "2026-03"}}
      }
    ]
  }
}
```

## Peek [#peek]

Peeking inspects queued messages **without consuming them** — pass `is_peek=true` on a `receive` call. The same messages remain available for a later consuming receive, which makes peek useful for monitoring queue depth or previewing work before committing to it.

<Tabs groupId="language" items="['curl','Go','JavaScript','Python']">
  <Tab value="curl">
    ```bash
    # Peek: inspect up to 10 messages without removing them
    curl -X POST "http://localhost:9090/ce/queue/receive?channel=work-queue&client_id=worker-1&max_messages=10&wait_timeout=3&is_peek=true"

    # A normal receive (is_peek omitted) consumes the same messages
    curl -X POST "http://localhost:9090/ce/queue/receive?channel=work-queue&client_id=worker-1&max_messages=10&wait_timeout=3"
    ```
  </Tab>

  <Tab value="Go">
    ```go
    func receiveOrPeek(base, channel, clientID string, isPeek bool, label string) int {
    	isPeekStr := "false"
    	if isPeek {
    		isPeekStr = "true"
    	}
    	url := fmt.Sprintf(
    		"%s/ce/queue/receive?channel=%s&client_id=%s&max_messages=10&wait_timeout=3&is_peek=%s",
    		base, channel, clientID, isPeekStr)
    	req, _ := http.NewRequest("POST", url, nil)
    	resp, err := http.DefaultClient.Do(req)
    	if err != nil {
    		log.Fatal("receive:", err)
    	}
    	defer resp.Body.Close()

    	var result CEResponse
    	_ = json.NewDecoder(resp.Body).Decode(&result)
    	var data QueueReceiveData
    	_ = json.Unmarshal(result.Data, &data)
    	fmt.Printf("[%s] messages_received=%d\n", label, data.MessagesReceived)
    	return data.MessagesReceived
    }

    // Peek twice (messages stay), then consume.
    receiveOrPeek(base, channel, "go-peek-client", true, "peek #1")
    receiveOrPeek(base, channel, "go-peek-client", true, "peek #2")
    receiveOrPeek(base, channel, "go-peek-client", false, "consume")
    ```
  </Tab>

  <Tab value="JavaScript">
    ```javascript
    async function receiveOrPeek(base, channel, clientId, isPeek, label) {
      const url = `${base}/ce/queue/receive?channel=${encodeURIComponent(channel)}&client_id=${clientId}&max_messages=10&wait_timeout=3&is_peek=${isPeek}`;
      const resp = await fetch(url, { method: 'POST' });
      const result = await resp.json();
      const n = result.data?.messages_received ?? 0;
      console.log(`[${label}] messages_received=${n}`);
      return n;
    }

    // Peek twice (messages stay), then consume.
    await receiveOrPeek(base, channel, 'js-peek-client', true, 'peek #1');
    await receiveOrPeek(base, channel, 'js-peek-client', true, 'peek #2');
    await receiveOrPeek(base, channel, 'js-peek-client', false, 'consume');
    ```
  </Tab>

  <Tab value="Python">
    ```python
    def receive_or_peek(base, channel, client_id, is_peek, label):
        resp = requests.post(
            f"{base}/ce/queue/receive",
            params={
                "channel": channel,
                "client_id": client_id,
                "max_messages": 10,
                "wait_timeout": 3,
                "is_peek": "true" if is_peek else "false",
            },
            timeout=10,
        )
        result = resp.json()
        n = result.get("data", {}).get("messages_received", 0)
        print(f"[{label}] messages_received={n}")
        return n


    # Peek twice (messages stay), then consume.
    receive_or_peek(base, channel, "python-peek-client", True, "peek #1")
    receive_or_peek(base, channel, "python-peek-client", True, "peek #2")
    receive_or_peek(base, channel, "python-peek-client", False, "consume")
    ```
  </Tab>
</Tabs>

<Callout type="info">
  The query parameter is spelled `is_peek` on the wire — setting `is_peek=true` performs a non-destructive peek (messages stay queued), while omitting it (or `is_peek=false`) performs a normal consuming receive. The internal Go struct field is named `IsPeak`, but that name never appears on the wire.
</Callout>

## Ack all [#ack-all]

`POST /ce/queue/ack_all` acknowledges (drains) **all pending messages** in a queue atomically, without receiving them one by one. This is a control operation — the body is ignored and `channel`, `client_id`, and `wait_timeout` come from query parameters. Use it to clear a backlog or reset a work channel.

<Tabs groupId="language" items="['curl','Go','Python']">
  <Tab value="curl">
    ```bash
    curl -X POST "http://localhost:9090/ce/queue/ack_all?channel=work-queue&client_id=worker-1&wait_timeout=10"
    ```
  </Tab>

  <Tab value="Go">
    ```go
    // Drain the queue atomically.
    ackURL := fmt.Sprintf("%s/ce/queue/ack_all?channel=%s&client_id=%s&wait_timeout=5",
    	base, channel, clientID)
    req, _ := http.NewRequest("POST", ackURL, nil)
    resp, err := http.DefaultClient.Do(req)
    if err != nil {
    	log.Fatal("ack_all:", err)
    }
    defer resp.Body.Close()

    var ackResult CEResponse
    _ = json.NewDecoder(resp.Body).Decode(&ackResult)
    fmt.Printf("ack_all: is_error=%v message=%s\n", ackResult.IsError, ackResult.Message)
    ```
  </Tab>

  <Tab value="Python">
    ```python
    # Drain the queue atomically.
    ack_resp = requests.post(
        f"{base}/ce/queue/ack_all",
        params={"channel": channel, "client_id": client_id, "wait_timeout": 5},
        timeout=10,
    )
    result = ack_resp.json()
    print(f"ack_all: is_error={result.get('is_error')} message={result.get('message')}")
    ```
  </Tab>
</Tabs>

### Ack-all parameters [#ack-all-parameters]

| Parameter      | Type   | Default      | Description                                                        |
| -------------- | ------ | ------------ | ------------------------------------------------------------------ |
| `channel`      | string | *(required)* | Queue channel name                                                 |
| `client_id`    | string | *(required)* | Client identifier (overridden by auth claims when auth is enabled) |
| `wait_timeout` | int    | `5`          | Wait timeout in seconds                                            |

A successful ack returns HTTP 200 with the result in `data`.

<Callout type="info">
  `send`, `receive`, and `ack_all` are synchronous `POST` endpoints and are subject to the connector's `TimeoutSeconds` (default 60s). A request that exceeds the timeout returns HTTP 504. See [Configuration](/connectors/cloudevents/concepts/configuration-model).
</Callout>

## Related [#related]

<Cards>
  <Card title="Endpoints reference" href="/connectors/cloudevents/reference/endpoints" description="Full CloudEvents endpoint table and status codes." />

  <Card title="Content modes" href="/connectors/cloudevents/how-to/content-modes" description="Structured vs binary CloudEvent encodings for queue sends." />

  <Card title="Authentication" href="/connectors/cloudevents/how-to/authentication" description="JWT Bearer auth and ClientID resolution for /ce/* endpoints." />
</Cards>
