Class: KubeMQ::QueuesClient

Inherits:
BaseClient show all
Defined in:
lib/kubemq/queues/client.rb

Overview

Client for KubeMQ message queues — guaranteed delivery with acknowledgement.

Provides two APIs: a stream API (primary, recommended) using persistent gRPC bidirectional streams for high throughput, and a simple API (secondary) using unary RPCs for low-volume use cases.

Inherits connection management and channel CRUD from BaseClient.

Examples:

Stream API — send and poll

client = KubeMQ::QueuesClient.new(address: "localhost:50000")

msg = KubeMQ::Queues::QueueMessage.new(channel: "tasks", body: "work-item")
client.send_queue_message_stream(msg)

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
client.close

See Also:

Instance Attribute Summary

Attributes inherited from BaseClient

#config, #transport

Instance Method Summary collapse

Methods inherited from BaseClient

#closed?, #create_channel, #delete_channel, #list_channels, #ping, #purge_queue_channel

Constructor Details

#initialize(address: nil, client_id: nil, auth_token: nil, config: nil, **options) ⇒ QueuesClient

Returns a new instance of QueuesClient.



34
35
36
37
38
39
40
41
# File 'lib/kubemq/queues/client.rb', line 34

def initialize(address: nil, client_id: nil, auth_token: nil, config: nil, **options)
  super
  @mutex = Mutex.new
  @closed = false
  @stream_mutex = Mutex.new
  @upstream_sender = nil
  @downstream_receiver = nil
end

Instance Method Details

#ack_all_queue_messages(channel:, wait_timeout_seconds: 1) ⇒ Integer

Note:

wait_timeout_seconds is in seconds.

Acknowledges all pending messages in a queue channel.

Parameters:

  • channel (String)

    queue channel to acknowledge

  • wait_timeout_seconds (Integer) (defaults to: 1)

    seconds to wait for the operation to complete (default: 1)

Returns:

  • (Integer)

    number of affected messages

Raises:



289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
# File 'lib/kubemq/queues/client.rb', line 289

def ack_all_queue_messages(channel:, wait_timeout_seconds: 1)
  Validator.validate_channel!(channel, allow_wildcards: false)
  ensure_connected!

  request = ::Kubemq::AckAllQueueMessagesRequest.new(
    RequestID: SecureRandom.uuid,
    ClientID: @config.client_id,
    Channel: channel,
    WaitTimeSeconds: wait_timeout_seconds
  )

  response = @transport.kubemq_client.ack_all_queue_messages(request)

  if response.IsError && !response.Error.empty?
    raise MessageError.new(response.Error, operation: 'ack_all_queue_messages', channel: channel)
  end

  response.AffectedMessages
end

#closevoid

This method returns an undefined value.

Closes the client, all open stream senders/receivers, and releases resources.

Closes any internal KubeMQ::Queues::UpstreamSender and KubeMQ::Queues::DownstreamReceiver before closing the transport. This method is idempotent.



350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
# File 'lib/kubemq/queues/client.rb', line 350

def close
  @mutex.synchronize do
    return if @closed

    @closed = true
  end
  @stream_mutex.synchronize do
    begin; @upstream_sender&.close; rescue StandardError; end
    @upstream_sender = nil
    begin; @downstream_receiver&.close; rescue StandardError; end
    @downstream_receiver = nil
  end
ensure
  super
end

#create_downstream_receiverQueues::DownstreamReceiver

Creates a new downstream receiver for polling queue messages over a persistent gRPC stream.

Returns:

Raises:

See Also:



68
69
70
71
# File 'lib/kubemq/queues/client.rb', line 68

def create_downstream_receiver
  ensure_connected!
  Queues::DownstreamReceiver.new(transport: @transport, client_id: @config.client_id)
end

#create_queues_channel(channel_name:) ⇒ Boolean

Creates a queues channel on the broker.

Parameters:

  • channel_name (String)

    name for the new channel

Returns:

  • (Boolean)

    true on success

Raises:

See Also:



318
319
320
# File 'lib/kubemq/queues/client.rb', line 318

def create_queues_channel(channel_name:)
  create_channel(channel_name: channel_name, channel_type: ChannelType::QUEUES)
end

#create_upstream_senderQueues::UpstreamSender

Creates a new upstream sender for publishing queue messages over a persistent gRPC stream.

Returns:

Raises:

See Also:



54
55
56
57
# File 'lib/kubemq/queues/client.rb', line 54

def create_upstream_sender
  ensure_connected!
  Queues::UpstreamSender.new(transport: @transport, client_id: @config.client_id)
end

#delete_queues_channel(channel_name:) ⇒ Boolean

Deletes a queues channel from the broker.

Parameters:

  • channel_name (String)

    name of the channel to delete

Returns:

  • (Boolean)

    true on success

Raises:

See Also:



329
330
331
# File 'lib/kubemq/queues/client.rb', line 329

def delete_queues_channel(channel_name:)
  delete_channel(channel_name: channel_name, channel_type: ChannelType::QUEUES)
end

#list_queues_channels(search: nil) ⇒ Array<ChannelInfo>

Lists queues channels, with optional name filtering.

Parameters:

  • search (String, nil) (defaults to: nil)

    substring filter for channel names

Returns:

  • (Array<ChannelInfo>)

    matching channels with metadata

Raises:

See Also:



339
340
341
# File 'lib/kubemq/queues/client.rb', line 339

def list_queues_channels(search: nil)
  list_channels(channel_type: ChannelType::QUEUES, search: search)
end

#poll(request) ⇒ Queues::QueuePollResponse

Polls for queue messages using the stream API.

Lazily creates an internal KubeMQ::Queues::DownstreamReceiver on first call. If the stream breaks, the receiver is reset and a StreamBrokenError is raised — retry the call to establish a new stream.

Examples:

request = KubeMQ::Queues::QueuePollRequest.new(
  channel: "tasks", max_items: 10, wait_timeout: 5
)
response = client.poll(request)
response.messages.each { |m| m.ack }

Parameters:

Returns:

Raises:

See Also:



126
127
128
129
130
131
132
133
134
# File 'lib/kubemq/queues/client.rb', line 126

def poll(request)
  receiver = @stream_mutex.synchronize do
    @downstream_receiver ||= create_downstream_receiver
  end
  receiver.poll(request)
rescue StreamBrokenError
  @stream_mutex.synchronize { @downstream_receiver = nil }
  raise
end

#receive_queue_messages(channel:, max_messages: 1, wait_timeout_seconds: 1, peek: false) ⇒ Array<Queues::QueueMessageReceived>

Note:

wait_timeout_seconds is in seconds.

Receives queue messages via a unary RPC call (simple API).

For higher throughput and transactional ack/nack/requeue, prefer the stream-based #poll method instead.

Examples:

messages = client.receive_queue_messages(
  channel: "tasks",
  max_messages: 5,
  wait_timeout_seconds: 10
)
messages.each { |m| puts m.body }

Parameters:

  • channel (String)

    queue channel to receive from

  • max_messages (Integer) (defaults to: 1)

    maximum number of messages to return (default: 1)

  • wait_timeout_seconds (Integer) (defaults to: 1)

    seconds to wait for messages before returning empty (default: 1)

  • peek (Boolean) (defaults to: false)

    if true, messages are not removed from the queue (default: false)

Returns:

Raises:

See Also:



241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
# File 'lib/kubemq/queues/client.rb', line 241

def receive_queue_messages(channel:, max_messages: 1, wait_timeout_seconds: 1, peek: false)
  Validator.validate_channel!(channel, allow_wildcards: false)
  Validator.validate_queue_receive!(max_messages, wait_timeout_seconds)
  ensure_connected!

  request = ::Kubemq::ReceiveQueueMessagesRequest.new(
    RequestID: SecureRandom.uuid,
    ClientID: @config.client_id,
    Channel: channel,
    MaxNumberOfMessages: max_messages,
    WaitTimeSeconds: wait_timeout_seconds,
    IsPeak: peek
  )

  response = @transport.kubemq_client.receive_queue_messages(request)

  if response.IsError && !response.Error.empty?
    raise MessageError.new(response.Error, operation: 'receive_queue_messages', channel: channel)
  end

  (response.Messages || []).map do |msg|
    hash = Transport::Converter.proto_to_queue_message_received(msg)
    attrs = (Queues::QueueMessageAttributes.new(**hash[:attributes]) if hash[:attributes])
    Queues::QueueMessageReceived.new(
      id: hash[:id],
      channel: hash[:channel],
      metadata: hash[:metadata],
      body: hash[:body],
      tags: hash[:tags],
      attributes: attrs
    )
  end
end

#send_queue_message(message) ⇒ Queues::QueueSendResult

Sends a single queue message via a unary RPC call.

For higher throughput, prefer #send_queue_message_stream or #create_upstream_sender instead.

Parameters:

Returns:

Raises:

See Also:



153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
# File 'lib/kubemq/queues/client.rb', line 153

def send_queue_message(message)
  Validator.validate_channel!(message.channel, allow_wildcards: false)
  ensure_connected!

  proto = Transport::Converter.queue_message_to_proto(message, @config.client_id)
  result = @transport.kubemq_client.send_queue_message(proto)

  Queues::QueueSendResult.new(
    id: result.MessageID,
    sent_at: result.SentAt,
    expiration_at: result.ExpirationAt,
    delayed_to: result.DelayedTo,
    error: result.IsError ? result.Error : nil
  )
end

#send_queue_message_stream(message) ⇒ Queues::QueueSendResult

Note:

Auto-creates an upstream sender on first call.

Sends a single queue message via the stream API.

Lazily creates an internal KubeMQ::Queues::UpstreamSender on first call. If the stream breaks, the sender is reset and a StreamBrokenError is raised — retry the call to establish a new stream.

Parameters:

Returns:

Raises:

See Also:



90
91
92
93
94
95
96
97
98
99
100
# File 'lib/kubemq/queues/client.rb', line 90

def send_queue_message_stream(message)
  Validator.validate_channel!(message.channel, allow_wildcards: false)
  sender = @stream_mutex.synchronize do
    @upstream_sender ||= create_upstream_sender
  end
  results = sender.publish([message])
  results.first
rescue StreamBrokenError
  @stream_mutex.synchronize { @upstream_sender = nil }
  raise
end

#send_queue_messages_batch(messages, batch_id: nil) ⇒ Array<Queues::QueueSendResult>

Sends multiple queue messages in a single batch via a unary RPC call.

Parameters:

  • messages (Array<Queues::QueueMessage>)

    messages to enqueue

  • batch_id (String, nil) (defaults to: nil)

    identifier for this batch (auto-generated UUID if nil)

Returns:

Raises:

See Also:



183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
# File 'lib/kubemq/queues/client.rb', line 183

def send_queue_messages_batch(messages, batch_id: nil)
  raise ValidationError, 'messages cannot be empty' if messages.nil? || messages.empty?

  ensure_connected!

  proto_messages = messages.map do |msg|
    Validator.validate_channel!(msg.channel, allow_wildcards: false)
    Transport::Converter.queue_message_to_proto(msg, @config.client_id)
  end

  batch_request = ::Kubemq::QueueMessagesBatchRequest.new(
    BatchID: batch_id || SecureRandom.uuid,
    Messages: proto_messages
  )

  response = @transport.kubemq_client.send_queue_messages_batch(batch_request)
  response.Results.map do |r|
    Queues::QueueSendResult.new(
      id: r.MessageID,
      sent_at: r.SentAt,
      expiration_at: r.ExpirationAt,
      delayed_to: r.DelayedTo,
      error: r.IsError ? r.Error : nil
    )
  end
end