KubeMQ
ConnectorsCloudEventsHow-to guides

Queues

Send, receive, peek, and ack durable FIFO queue messages over CloudEvents HTTP with the KubeMQ CloudEvents connector.

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

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:

OperationEndpointBodyPurpose
SendPOST /ce/queue/sendCloudEvent (structured or binary)Enqueue one message
ReceivePOST /ce/queue/receive(ignored) — all params via query stringPull (and consume) messages
Ack allPOST /ce/queue/ack_all(ignored) — all params via query stringDrain 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

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

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

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.

# 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"
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")}");
    }
}
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)
	}
}
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"));
            }
        }
    }
}
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);
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()
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
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(())
}

Receive parameters

ParameterTypeDefaultDescription
channelstring(required)Queue channel name
client_idstring(required)Client identifier (overridden by auth claims when auth is enabled)
max_messagesint1Maximum number of messages to receive (1–1000)
wait_timeoutint5Wait timeout in seconds
is_peekboolfalseIf 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:

{
  "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

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.

# 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"
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")
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');
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")

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.

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.

curl -X POST "http://localhost:9090/ce/queue/ack_all?channel=work-queue&client_id=worker-1&wait_timeout=10"
// 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)
# 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')}")

Ack-all parameters

ParameterTypeDefaultDescription
channelstring(required)Queue channel name
client_idstring(required)Client identifier (overridden by auth claims when auth is enabled)
wait_timeoutint5Wait timeout in seconds

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

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.

Was this page helpful?

On this page