Queues
KubeMQ Ruby SDK API reference for guaranteed message delivery with queues.
QueueMessage
Outbound queue message for point-to-point messaging. Optionally attach a QueueMessagePolicy for expiration, delay, and dead-letter handling.
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
| 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
Delivery policy for queue messages — expiration, delay, and dead-letter settings.
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
Configuration for polling messages from a queue channel.
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
Stream API (Primary)
send_queue_message_stream(message)
Sends a queue message via the persistent gRPC stream. Auto-creates an upstream sender on first call.
result = client.send_queue_message_stream(msg)
puts "Sent: #{result.id}"Returns: QueueSendResult
poll(request)
Polls for queue messages using the stream API. Returns a response with transactional ack/nack/requeue support.
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
endReturns: QueuePollResponse
create_upstream_sender
Creates a new upstream sender for publishing via persistent gRPC stream.
Returns: Queues::UpstreamSender
create_downstream_receiver
Creates a new downstream receiver for polling via persistent gRPC stream.
Returns: Queues::DownstreamReceiver
Simple API (Secondary)
send_queue_message(message)
Sends a single queue message via a unary RPC call.
send_queue_messages_batch(messages, batch_id:)
Sends multiple queue messages in a single batch.
results = client.send_queue_messages_batch(messages)Returns: Array<QueueSendResult>
receive_queue_messages(channel:, max_messages:, wait_timeout_seconds:, peek:)
Receives queue messages via a unary RPC call.
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:)
Acknowledges all pending messages in a queue channel.
affected = client.ack_all_queue_messages(channel: "tasks", wait_timeout_seconds: 5)Returns: Integer — number of affected messages
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
| 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
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 |
Was this page helpful?