KubeMQ
Client SDKsRubyReference

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

AttributeTypeDefaultDescription
idStringAuto-generated UUIDUnique message identifier
channelStringRequiredTarget queue channel
metadataStringnilArbitrary metadata
bodyStringnilMessage payload
tagsHash{String => String}{}Key-value tags
policyQueueMessagePolicynilDelivery 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"
)
AttributeTypeDefaultDescription
expiration_secondsInteger0Message TTL in seconds; 0 means no expiration
delay_secondsInteger0Delay before message becomes visible
max_receive_countInteger0Max delivery attempts before dead-lettering; 0 means unlimited
max_receive_queueString""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
)
AttributeTypeDefaultDescription
channelStringRequiredQueue channel to poll from
max_itemsInteger1Maximum messages to receive
wait_timeoutInteger1Wait timeout in seconds
auto_ackBooleanfalseAuto-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
end

Returns: 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.

AttributeTypeDescription
idStringMessage identifier
channelStringSource channel
metadataStringMetadata string
bodyStringMessage payload
tagsHashKey-value tags
attributesQueueMessageAttributesDelivery metadata

Transaction Methods

MethodDescription
ackAcknowledge the message (removes from queue)
rejectNegative-acknowledge (returns to queue)
requeue(channel:)Move the message to a different channel

QueuePollResponse

Response from a poll operation with batch-level actions.

MethodDescription
messagesArray<QueueMessageReceived> — received messages
error?Whether the poll encountered an error
errorError message string
ack_allAcknowledge all messages in the batch
nack_allNegative-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?

On this page