# Queues (/sdks/ruby/reference/queues)



## QueueMessage [#queuemessage]

Outbound queue message for point-to-point messaging. Optionally attach a `QueueMessagePolicy` for expiration, delay, and dead-letter handling.

```ruby
msg = KubeMQ::Queues::QueueMessage.new(
  channel: "tasks.process",
  body: '{"task_id": 1}',
  metadata: "task",
  tags: { "priority" => "high" },
  policy: KubeMQ::Queues::QueueMessagePolicy.new(
    expiration_seconds: 300,
    max_receive_count: 3,
    max_receive_queue: "tasks.dlq"
  )
)
```

### Attributes [#attributes]

| Attribute  | Type                     | Default             | Description               |
| ---------- | ------------------------ | ------------------- | ------------------------- |
| `id`       | `String`                 | Auto-generated UUID | Unique message identifier |
| `channel`  | `String`                 | Required            | Target queue channel      |
| `metadata` | `String`                 | `nil`               | Arbitrary metadata        |
| `body`     | `String`                 | `nil`               | Message payload           |
| `tags`     | `Hash{String => String}` | `{}`                | Key-value tags            |
| `policy`   | `QueueMessagePolicy`     | `nil`               | Delivery policy           |

## QueueMessagePolicy [#queuemessagepolicy]

Delivery policy for queue messages — expiration, delay, and dead-letter settings.

```ruby
policy = KubeMQ::Queues::QueueMessagePolicy.new(
  expiration_seconds: 3600,
  delay_seconds: 30,
  max_receive_count: 3,
  max_receive_queue: "orders.dlq"
)
```

| Attribute            | Type      | Default | Description                                                      |
| -------------------- | --------- | ------- | ---------------------------------------------------------------- |
| `expiration_seconds` | `Integer` | `0`     | Message TTL in seconds; `0` means no expiration                  |
| `delay_seconds`      | `Integer` | `0`     | Delay before message becomes visible                             |
| `max_receive_count`  | `Integer` | `0`     | Max delivery attempts before dead-lettering; `0` means unlimited |
| `max_receive_queue`  | `String`  | `""`    | Dead-letter queue channel name                                   |

## QueuePollRequest [#queuepollrequest]

Configuration for polling messages from a queue channel.

```ruby
request = KubeMQ::Queues::QueuePollRequest.new(
  channel: "tasks",
  max_items: 10,
  wait_timeout: 5,
  auto_ack: false
)
```

| Attribute      | Type      | Default  | Description                 |
| -------------- | --------- | -------- | --------------------------- |
| `channel`      | `String`  | Required | Queue channel to poll from  |
| `max_items`    | `Integer` | `1`      | Maximum messages to receive |
| `wait_timeout` | `Integer` | `1`      | Wait timeout in seconds     |
| `auto_ack`     | `Boolean` | `false`  | Auto-acknowledge on receipt |

## QueuesClient Methods [#queuesclient-methods]

### Stream API (Primary) [#stream-api-primary]

#### `send_queue_message_stream(message)` [#send_queue_message_streammessage]

Sends a queue message via the persistent gRPC stream. Auto-creates an upstream sender on first call.

```ruby
result = client.send_queue_message_stream(msg)
puts "Sent: #{result.id}"
```

**Returns:** `QueueSendResult`

#### `poll(request)` [#pollrequest]

Polls for queue messages using the stream API. Returns a response with transactional ack/nack/requeue support.

```ruby
request = KubeMQ::Queues::QueuePollRequest.new(channel: "tasks", max_items: 10, wait_timeout: 5)
response = client.poll(request)
response.messages.each do |m|
  process(m.body)
  m.ack
end
```

**Returns:** `QueuePollResponse`

#### `create_upstream_sender` [#create_upstream_sender]

Creates a new upstream sender for publishing via persistent gRPC stream.

**Returns:** `Queues::UpstreamSender`

#### `create_downstream_receiver` [#create_downstream_receiver]

Creates a new downstream receiver for polling via persistent gRPC stream.

**Returns:** `Queues::DownstreamReceiver`

### Simple API (Secondary) [#simple-api-secondary]

#### `send_queue_message(message)` [#send_queue_messagemessage]

Sends a single queue message via a unary RPC call.

#### `send_queue_messages_batch(messages, batch_id:)` [#send_queue_messages_batchmessages-batch_id]

Sends multiple queue messages in a single batch.

```ruby
results = client.send_queue_messages_batch(messages)
```

**Returns:** `Array<QueueSendResult>`

#### `receive_queue_messages(channel:, max_messages:, wait_timeout_seconds:, peek:)` [#receive_queue_messageschannel-max_messages-wait_timeout_seconds-peek]

Receives queue messages via a unary RPC call.

```ruby
messages = client.receive_queue_messages(
  channel: "tasks",
  max_messages: 5,
  wait_timeout_seconds: 10,
  peek: false
)
```

**Returns:** `Array<QueueMessageReceived>`

#### `ack_all_queue_messages(channel:, wait_timeout_seconds:)` [#ack_all_queue_messageschannel-wait_timeout_seconds]

Acknowledges all pending messages in a queue channel.

```ruby
affected = client.ack_all_queue_messages(channel: "tasks", wait_timeout_seconds: 5)
```

**Returns:** `Integer` — number of affected messages

## QueueMessageReceived [#queuemessagereceived]

Inbound queue message with transaction methods.

| Attribute    | Type                     | Description        |
| ------------ | ------------------------ | ------------------ |
| `id`         | `String`                 | Message identifier |
| `channel`    | `String`                 | Source channel     |
| `metadata`   | `String`                 | Metadata string    |
| `body`       | `String`                 | Message payload    |
| `tags`       | `Hash`                   | Key-value tags     |
| `attributes` | `QueueMessageAttributes` | Delivery metadata  |

### Transaction Methods [#transaction-methods]

| Method              | Description                                  |
| ------------------- | -------------------------------------------- |
| `ack`               | Acknowledge the message (removes from queue) |
| `reject`            | Negative-acknowledge (returns to queue)      |
| `requeue(channel:)` | Move the message to a different channel      |

## QueuePollResponse [#queuepollresponse]

Response from a poll operation with batch-level actions.

| Method                       | Description                                       |
| ---------------------------- | ------------------------------------------------- |
| `messages`                   | `Array<QueueMessageReceived>` — received messages |
| `error?`                     | Whether the poll encountered an error             |
| `error`                      | Error message string                              |
| `ack_all`                    | Acknowledge all messages in the batch             |
| `nack_all`                   | Negative-acknowledge all messages                 |
| `requeue_all(channel:)`      | Requeue all messages to a channel                 |
| `ack_range(sequence_range:)` | Acknowledge specific messages by sequence         |
