Class: KubeMQ::QueuesClient
- Inherits:
-
BaseClient
- Object
- BaseClient
- KubeMQ::QueuesClient
- 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.
Instance Attribute Summary
Attributes inherited from BaseClient
Instance Method Summary collapse
-
#ack_all_queue_messages(channel:, wait_timeout_seconds: 1) ⇒ Integer
Acknowledges all pending messages in a queue channel.
-
#close ⇒ void
Closes the client, all open stream senders/receivers, and releases resources.
-
#create_downstream_receiver ⇒ Queues::DownstreamReceiver
Creates a new downstream receiver for polling queue messages over a persistent gRPC stream.
-
#create_queues_channel(channel_name:) ⇒ Boolean
Creates a queues channel on the broker.
-
#create_upstream_sender ⇒ Queues::UpstreamSender
Creates a new upstream sender for publishing queue messages over a persistent gRPC stream.
-
#delete_queues_channel(channel_name:) ⇒ Boolean
Deletes a queues channel from the broker.
-
#initialize(address: nil, client_id: nil, auth_token: nil, config: nil, **options) ⇒ QueuesClient
constructor
A new instance of QueuesClient.
-
#list_queues_channels(search: nil) ⇒ Array<ChannelInfo>
Lists queues channels, with optional name filtering.
-
#poll(request) ⇒ Queues::QueuePollResponse
Polls for queue messages using the stream API.
-
#receive_queue_messages(channel:, max_messages: 1, wait_timeout_seconds: 1, peek: false) ⇒ Array<Queues::QueueMessageReceived>
Receives queue messages via a unary RPC call (simple API).
-
#send_queue_message(message) ⇒ Queues::QueueSendResult
Sends a single queue message via a unary RPC call.
-
#send_queue_message_stream(message) ⇒ Queues::QueueSendResult
Sends a single queue message via the stream API.
-
#send_queue_messages_batch(messages, batch_id: nil) ⇒ Array<Queues::QueueSendResult>
Sends multiple queue messages in a single batch via a unary RPC call.
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, **) 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
wait_timeout_seconds is in seconds.
Acknowledges all pending messages in a queue channel.
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 (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.(request) if response.IsError && !response.Error.empty? raise MessageError.new(response.Error, operation: 'ack_all_queue_messages', channel: channel) end response.AffectedMessages end |
#close ⇒ void
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_receiver ⇒ Queues::DownstreamReceiver
Creates a new downstream receiver for polling queue messages over a persistent gRPC stream.
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.
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_sender ⇒ Queues::UpstreamSender
Creates a new upstream sender for publishing queue messages over a persistent gRPC stream.
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.
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.
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.
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>
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.
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 (channel:, max_messages: 1, wait_timeout_seconds: 1, peek: false) Validator.validate_channel!(channel, allow_wildcards: false) Validator.validate_queue_receive!(, wait_timeout_seconds) ensure_connected! request = ::Kubemq::ReceiveQueueMessagesRequest.new( RequestID: SecureRandom.uuid, ClientID: @config.client_id, Channel: channel, MaxNumberOfMessages: , WaitTimeSeconds: wait_timeout_seconds, IsPeak: peek ) response = @transport.kubemq_client.(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.(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.
153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 |
# File 'lib/kubemq/queues/client.rb', line 153 def () Validator.validate_channel!(.channel, allow_wildcards: false) ensure_connected! proto = Transport::Converter.(, @config.client_id) result = @transport.kubemq_client.(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
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.
90 91 92 93 94 95 96 97 98 99 100 |
# File 'lib/kubemq/queues/client.rb', line 90 def () Validator.validate_channel!(.channel, allow_wildcards: false) sender = @stream_mutex.synchronize do @upstream_sender ||= create_upstream_sender end results = sender.publish([]) 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.
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 (, batch_id: nil) raise ValidationError, 'messages cannot be empty' if .nil? || .empty? ensure_connected! = .map do |msg| Validator.validate_channel!(msg.channel, allow_wildcards: false) Transport::Converter.(msg, @config.client_id) end batch_request = ::Kubemq::QueueMessagesBatchRequest.new( BatchID: batch_id || SecureRandom.uuid, Messages: ) response = @transport.kubemq_client.(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 |