# Queue Reference (/learn/queues/reference)



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

## Message Structure [#message-structure]

### Send Request [#send-request]

<TypeTable
  type="{
  messageId: { type: 'string', description: 'Unique identifier. Auto-generated UUID if empty.', default: 'Auto-generated' },
  channel: { type: 'string', description: 'Target queue channel name. Must follow channel naming rules.', required: true },
  clientId: { type: 'string', description: 'Identifier of the sending client.', default: 'Server-assigned' },
  metadata: { type: 'string', description: 'Text metadata. At least one of body or metadata is required.', default: '&#x22;&#x22;' },
  body: { type: 'bytes', description: 'Message payload. At least one of body or metadata is required.', default: 'null' },
  tags: { type: 'map<string, string>', description: 'Key-value pairs for filtering and routing.', default: '{}' },
}"
/>

### Message Policy Fields [#message-policy-fields]

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

<TypeTable
  type="{
  expirationSeconds: { type: 'int', description: 'Seconds after which the message expires and is discarded. 0 = no expiration.', default: '0' },
  delaySeconds: { type: 'int', description: 'Seconds before the message becomes available for consumption. 0 = immediate.', default: '0' },
  maxReceiveCount: { type: 'int', description: 'Maximum delivery attempts before DLQ routing. 0 = server default (1024).', default: '0' },
  maxReceiveQueue: { type: 'string', description: 'Dead letter queue channel for failed messages. Empty = discard on exceed.', default: '&#x22;&#x22;' },
}"
/>

### Send Response [#send-response]

<TypeTable
  type="{
  messageId: { type: 'string', description: 'Unique identifier assigned to the sent message.' },
  sentAt: { type: 'int (epoch ms)', description: 'Timestamp when the message was accepted by the server.' },
  expiresAt: { type: 'int (epoch ms)', description: 'When the message will expire. 0 = no expiration.' },
  delayedTo: { type: 'int (epoch ms)', description: 'When the message becomes available. 0 = immediate.' },
  isError: { type: 'boolean', description: 'Whether an error occurred during send.' },
  error: { type: 'string', description: 'Error description. Empty on success.' },
}"
/>

### Receive Response (per message) [#receive-response-per-message]

<TypeTable
  type="{
  messageId: { type: 'string', description: 'Unique identifier of the message.' },
  clientId: { type: 'string', description: 'Identifier of the client that sent the message.' },
  metadata: { type: 'string', description: 'Text metadata from the sender.' },
  body: { type: 'bytes', description: 'Message payload from the sender.' },
  timestamp: { type: 'int (epoch ms)', description: 'When the message was originally sent.' },
  sequence: { type: 'int', description: 'Sequential position within the queue.' },
  tags: { type: 'map<string, string>', description: 'Key-value pairs from the sender.' },
  receiveCount: { type: 'int', description: 'Number of times this message has been delivered.' },
  reRoutedFrom: { type: 'string', description: 'Original channel if rerouted from a dead letter queue.' },
  expirationAt: { type: 'int (epoch ms)', description: 'When the message expires. 0 = no expiration.' },
  delayedTo: { type: 'int (epoch ms)', description: 'Until when the message was delayed. 0 = no delay.' },
}"
/>

## Receive Request [#receive-request]

<TypeTable
  type="{
  channel: { type: 'string', description: 'Queue channel to receive from.', required: true },
  maxMessages: { type: 'int', description: 'Maximum messages to receive (1-1024).', required: true },
  waitTimeoutSeconds: { type: 'int', description: 'Seconds to wait for messages before returning.', default: '1' },
  isPeek: { type: 'boolean', description: 'When true, inspect messages without consuming them.', default: 'false' },
  autoAck: { type: 'boolean', description: 'When true, messages are acknowledged on receipt.', default: 'false' },
  visibilitySeconds: { type: 'int', description: 'Seconds the message stays hidden during processing.', default: '60' },
}"
/>

## Channel Name Rules [#channel-name-rules]

| Rule            | Constraint                | Error Code |
| --------------- | ------------------------- | ---------- |
| Required        | Cannot be empty           | 120        |
| No trailing dot | Cannot end with `.`       | 119        |
| No whitespace   | Cannot contain spaces     | 108        |
| No wildcards    | Cannot 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 [#server-configuration]

<Accordions>
  <Accordion title="Queue Policy Limits">
    <TypeTable
      type="{
  MaxNumberOfMessages: { type: 'int', description: 'Maximum messages per receive request.', default: '1024' },
  MaxWaitTimeoutSeconds: { type: 'int', description: 'Maximum wait timeout for receive operations.', default: '3600 (1 hour)' },
  MaxExpirationSeconds: { type: 'int', description: 'Maximum message expiration time.', default: '43200 (12 hours)' },
  MaxDelaySeconds: { type: 'int', description: 'Maximum message delay time.', default: '43200 (12 hours)' },
  MaxReceiveCount: { type: 'int', description: 'Maximum delivery attempts before dead letter routing.', default: '1024' },
  MaxVisibilitySeconds: { type: 'int', description: 'Maximum visibility timeout during processing.', default: '43200 (12 hours)' },
  DefaultVisibilitySeconds: { type: 'int', description: 'Default visibility timeout when not specified.', default: '60 (1 minute)' },
  DefaultWaitTimeoutSeconds: { type: 'int', description: 'Default wait timeout when not specified.', default: '1' },
  MaxInflight: { type: 'int', description: 'Maximum unacknowledged messages in-flight simultaneously.', default: '2048' },
}"
    />
  </Accordion>

  <Accordion title="Body Size Limits">
    | Transport | Default Limit | Setting                     |
    | --------- | ------------- | --------------------------- |
    | gRPC      | \~1 GB        | `Connectors.Grpc.BodyLimit` |
    | REST      | 100 MB        | `Connectors.Rest.BodyLimit` |
  </Accordion>
</Accordions>

## Message Lifecycle [#message-lifecycle]

<Mermaid
  chart="stateDiagram-v2
    [*] --> Delayed: Send with delay > 0
    [*] --> Available: Send (no delay)
    Delayed --> Available: Delay expires
    Available --> Processing: Consumer receives
    Processing --> Removed: ack()
    Processing --> Available: nack()
    Processing --> Available: Visibility timeout expires
    Processing --> DeadLetter: Receive count exceeds max
    Available --> Expired: Expiration time reached
    Expired --> [*]: Discarded on next receive
    Removed --> [*]
    DeadLetter --> [*]: Moved to DLQ channel"
/>

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

## Dead Letter Queue Behavior [#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 [#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 [#error-codes]

### Queue-Specific Errors [#queue-specific-errors]

| Code | Error                                | Cause                                             |
| ---- | ------------------------------------ | ------------------------------------------------- |
| 120  | Invalid queue name                   | Empty channel or invalid characters               |
| 122  | Max messages exceeded config limit   | Receive request exceeds `MaxNumberOfMessages`     |
| 125  | Ack sequence cannot be 0 or negative | Invalid sequence in stream ack                    |
| 127  | Invalid visibility time              | Visibility timeout exceeds `MaxVisibilitySeconds` |
| 133  | Invalid expiration seconds           | Expiration exceeds `MaxExpirationSeconds`         |
| 134  | Invalid max receive count            | Max receive count exceeds server config           |
| 135  | Invalid delay seconds                | Delay exceeds `MaxDelaySeconds`                   |
| 136  | Invalid wait timeout                 | Wait timeout exceeds `MaxWaitTimeoutSeconds`      |
| 139  | Invalid max items                    | Value must be at least 1                          |
| 140  | Invalid wait time                    | Value must be at least 1 second                   |

### General Errors [#general-errors]

| Code | Error                                  | Applies To                     |
| ---- | -------------------------------------- | ------------------------------ |
| 101  | Invalid clientID, cannot be empty      | All patterns                   |
| 107  | Invalid channel, no wildcards allowed  | All patterns                   |
| 108  | Invalid channel, no whitespace allowed | All patterns                   |
| 110  | Invalid message, cannot be empty       | Body and metadata both missing |
| 119  | Invalid channel, cannot end with dot   | All patterns                   |

## SDK Quick Reference [#sdk-quick-reference]

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

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

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

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

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

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

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

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

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

  <Tab value="Elixir">
    ```elixir
    # {: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")
    ```
  </Tab>
</Tabs>

## Related Pages [#related-pages]

* [Queues Overview](/learn/queues) — Introduction and key features
* [Getting Started](/learn/queues/getting-started) — Step-by-step first message guide
* [Ack, Nack & Requeue](/learn/queues/tutorials/ack-nack-requeue) — Settlement options
* [Dead Letter Queue](/learn/queues/tutorials/dead-letter-queue) — DLQ routing and monitoring
* [Delayed Messages](/learn/queues/tutorials/delayed-messages) — Schedule future delivery
* [Visibility Timeout](/learn/queues/how-to/visibility-timeout) — Processing time configuration
* [Batch Operations](/learn/queues/tutorials/batch-operations) — High-throughput send and receive
