KubeMQ
LearnQueues

Queue Reference

Complete reference for KubeMQ queue message structure, policy fields, server configuration, and error codes.

This reference documents every aspect of the KubeMQ Queues messaging pattern.

Message Structure

Send Request

Prop

Type

Message Policy Fields

Policy fields control delivery behavior and are set per message at send time.

Prop

Type

Send Response

Prop

Type

Receive Response (per message)

Prop

Type

Receive Request

Prop

Type

Channel Name Rules

RuleConstraintError Code
RequiredCannot be empty120
No trailing dotCannot end with .119
No whitespaceCannot contain spaces108
No wildcardsCannot contain * or >107

Valid channel name pattern: ^[^\s*>]+[^.]$

Examples:

  • orders — valid
  • payments.us-east — valid
  • orders processing — invalid (contains space)
  • orders. — invalid (trailing dot)
  • orders.* — invalid (contains wildcard)

Server Configuration

Message Lifecycle

A queue message moves through delay, availability, processing, and settlement — ending in removal, redelivery, expiration, or the dead letter queue.

Dead Letter Queue Behavior

When a message exceeds the maximum receive count:

  1. If maxReceiveQueue is set, the message is copied to the DLQ channel with:
    • reRouted set to true
    • reRoutedFromQueue set to the original channel name
    • receiveCount reset to 0
    • Policy reset to server defaults
  2. If maxReceiveQueue is not set, the message is silently discarded
  3. The original message is acknowledged (removed from the source queue)

Delayed Message Behavior

Messages with delaySeconds > 0 follow this path:

  1. Message is published to an internal delay channel (_QUEUE_DELAY_)
  2. A background processor checks for expired delays every 500ms
  3. When the delay elapses, the message is moved to the target queue
  4. The delaySeconds and delayedTo attributes are cleared before delivery
  5. If both delay and expiration are set, the expiration clock starts after the delay

Error Codes

Queue-Specific Errors

CodeErrorCause
120Invalid queue nameEmpty channel or invalid characters
122Max messages exceeded config limitReceive request exceeds MaxNumberOfMessages
125Ack sequence cannot be 0 or negativeInvalid sequence in stream ack
127Invalid visibility timeVisibility timeout exceeds MaxVisibilitySeconds
133Invalid expiration secondsExpiration exceeds MaxExpirationSeconds
134Invalid max receive countMax receive count exceeds server config
135Invalid delay secondsDelay exceeds MaxDelaySeconds
136Invalid wait timeoutWait timeout exceeds MaxWaitTimeoutSeconds
139Invalid max itemsValue must be at least 1
140Invalid wait timeValue must be at least 1 second

General Errors

CodeErrorApplies To
101Invalid clientID, cannot be emptyAll patterns
107Invalid channel, no wildcards allowedAll patterns
108Invalid channel, no whitespace allowedAll patterns
110Invalid message, cannot be emptyBody and metadata both missing
119Invalid channel, cannot end with dotAll patterns

SDK Quick Reference

// go get github.com/kubemq-io/kubemq-go/v2

msg := kubemq.NewQueueMessage().
    SetChannel("orders").
    SetBody([]byte("payload")).
    SetMetadata("metadata").
    SetTags(map[string]string{"key": "value"}).
    SetDelaySeconds(10).
    SetExpirationSeconds(300).
    SetMaxReceiveCount(3).
    SetMaxReceiveQueue("orders.dlq")

result, err := client.SendQueueMessage(ctx, msg)
results, err := client.SendQueueMessages(ctx, msgs)

resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
    Channel: "orders", MaxItems: 10, WaitTimeoutSeconds: 5,
    AutoAck: false, VisibilitySeconds: 120, IsPeek: false,
})
resp.AckAll()
resp.NAckAll()
resp.ReQueueAll("other-channel")
# pip install kubemq

from kubemq.queues import Client as QueuesClient
from kubemq import QueueMessage

result = client.send_queue_message(QueueMessage(
    channel="orders", body=b"payload", metadata="metadata",
    tags={"key": "value"}, delay_in_seconds=10,
    expiration_in_seconds=300, max_receive_count=3,
    max_receive_queue="orders.dlq",
))

response = client.receive_queue_messages(
    channel="orders", max_messages=10, wait_timeout_in_seconds=5,
    is_peek=False, visibility_seconds=120,
)

msg.ack()
msg.nack()
msg.requeue("other-channel")
// npm install kubemq-js

import { KubeMQClient, createQueueMessage } from 'kubemq-js';

const result = await client.sendQueueMessage(createQueueMessage({
  channel: 'orders', body: 'payload', metadata: 'metadata',
  tags: { key: 'value' },
  policy: { delaySeconds: 10, expirationSeconds: 300,
    maxReceiveCount: 3, maxReceiveQueue: 'orders.dlq' },
}));

const msgs = await client.receiveQueueMessages({
  channel: 'orders', maxMessages: 10, waitTimeoutSeconds: 5,
  isPeek: false, visibilitySeconds: 120,
});

await msg.ack();
await msg.nack();
await msg.requeue('other-channel');
// io.kubemq.sdk:kubemq-sdk-Java:2.1.1

QueueMessage msg = QueueMessage.builder()
    .channel("orders").body("payload".getBytes())
    .metadata("metadata").tags(Map.of("key", "value"))
    .delaySeconds(10).expirationSeconds(300)
    .maxReceiveCount(3).maxReceiveQueue("orders.dlq").build();

SendQueueMessageResult result = client.sendQueueMessage(msg);

ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
    ReceiveQueueMessagesRequest.builder()
        .channel("orders").maxMessages(10).waitTimeoutSeconds(5)
        .isPeek(false).visibilitySeconds(120).build());

msg.ack();
msg.nack();
msg.requeue("other-channel");
// dotnet add package KubeMQ.SDK.CSharp

var result = await client.SendQueueMessageAsync(new QueueMessage {
    Channel = "orders", Body = Encoding.UTF8.GetBytes("payload"),
    Metadata = "metadata", Tags = new() { ["key"] = "value" },
    DelaySeconds = 10, ExpirationSeconds = 300,
    MaxReceiveCount = 3, MaxReceiveQueue = "orders.dlq"
});

var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest {
    Channel = "orders", MaxMessages = 10, WaitTimeoutSeconds = 5,
    IsPeek = false, VisibilitySeconds = 120,
});

await msg.AckAsync();
await msg.NAckAsync();
await msg.ReQueueAsync("other-channel");
// io.kubemq.sdk:kubemq-sdk-kotlin:2.1.0

val result = client.sendQueueMessage(QueueMessage(
    channel = "orders", body = "payload".toByteArray(),
    metadata = "metadata", tags = mapOf("key" to "value"),
    delaySeconds = 10, expirationSeconds = 300,
    maxReceiveCount = 3, maxReceiveQueue = "orders.dlq"
))

val response = client.receiveQueueMessages(
    channel = "orders", maxMessages = 10, waitTimeoutSeconds = 5,
    isPeek = false, visibilitySeconds = 120)

msg.ack()
msg.nack()
msg.requeue("other-channel")
// vcpkg install kubemq

kubemq::QueueMessage msg;
msg.channel = "orders";
msg.body = "payload";
msg.metadata = "metadata";
msg.tags = {{"key", "value"}};
msg.delaySeconds = 10;
msg.expirationSeconds = 300;
msg.maxReceiveCount = 3;
msg.maxReceiveQueue = "orders.dlq";

auto result = client.sendQueueMessage(msg);
auto response = client.receiveQueueMessages("orders", 10, 5);

msg.ack();
msg.nack();
msg.requeue("other-channel");
// cargo add kubemq

use kubemq::prelude::*;
use kubemq::{PollRequest, QueueMessageBuilder};

let msg = QueueMessageBuilder::new()
    .channel("orders")
    .body(b"payload".to_vec())
    .metadata("metadata")
    .tags([("key", "value")])
    .delay_seconds(10)
    .expiration_seconds(300)
    .max_receive_count(3)
    .max_receive_queue("orders.dlq")
    .build();

let result = client.send_queue_message(msg).await?;

// receive_queue_messages(channel, max_items, wait_timeout_seconds, is_peek)
let messages = client.receive_queue_messages("orders", 10, 5, false).await?;

// Stream poll for explicit settlement (ack / nack / re-queue)
let mut receiver = client.new_queue_downstream_receiver().await?;
let response = receiver.poll(PollRequest {
    channel: "orders".to_string(),
    max_items: 10, wait_timeout_seconds: 5, auto_ack: false,
}).await?;

response.ack_all().await?;
response.nack_all().await?;
for m in &response.messages { m.re_queue("other-channel").await?; }
# gem install kubemq

require 'kubemq'

policy = KubeMQ::Queues::QueueMessagePolicy.new(
  delay_seconds: 10, expiration_seconds: 300,
  max_receive_count: 3, max_receive_queue: 'orders.dlq'
)
msg = KubeMQ::Queues::QueueMessage.new(
  channel: 'orders', body: 'payload', metadata: 'metadata',
  tags: { 'key' => 'value' }, policy: policy
)
result = client.send_queue_message(msg)

messages = client.receive_queue_messages(
  channel: 'orders', max_messages: 10, wait_timeout_seconds: 5
)

# Stream poll for explicit settlement (ack / nack / requeue)
receiver = client.create_downstream_receiver
request  = KubeMQ::Queues::QueuePollRequest.new(
  channel: 'orders', max_items: 10, wait_timeout: 5
)
response = receiver.poll(request)
response.messages.each(&:ack)
response.messages.each(&:nack)
response.requeue_all(channel: 'other-channel')
# {:kubemq, "~> 1.0"} in mix.exs

msg = KubeMQ.QueueMessage.new(
  channel: "orders", body: "payload", metadata: "metadata",
  tags: %{"key" => "value"},
  policy: KubeMQ.QueuePolicy.new(
    delay_seconds: 10, expiration_seconds: 300,
    max_receive_count: 3, max_receive_queue: "orders.dlq"
  )
)
{:ok, result} = KubeMQ.Client.send_queue_message(client, msg)

{:ok, received} = KubeMQ.Client.receive_queue_messages(client, "orders",
  max_messages: 10, wait_timeout: 5_000)

# Stream poll for explicit settlement (ack / nack / requeue)
{:ok, poll} = KubeMQ.Client.poll_queue(client,
  channel: "orders", max_items: 10, wait_timeout: 5_000)

KubeMQ.PollResponse.ack_all(poll)
KubeMQ.PollResponse.nack_all(poll)
KubeMQ.PollResponse.requeue_all(poll, "other-channel")

Was this page helpful?

On this page