Class: KubeMQ::Queues::QueuePollResponse
- Inherits:
-
Object
- Object
- KubeMQ::Queues::QueuePollResponse
- Defined in:
- lib/kubemq/queues/queue_poll_response.rb
Overview
Response from polling a queue channel via KubeMQ::QueuesClient#poll.
Contains the received #messages and provides batch transaction actions (#ack_all, #nack_all, #requeue_all) as well as range-based variants for partial acknowledgement.
Instance Attribute Summary collapse
-
#active_offsets ⇒ Array<Integer>
readonly
Currently active message offsets.
-
#error ⇒ String?
readonly
Error description if the poll failed.
-
#messages ⇒ Array<QueueMessageReceived>
readonly
Received messages.
-
#transaction_complete ⇒ Object
readonly
Returns the value of attribute transaction_complete.
-
#transaction_id ⇒ String
readonly
Broker-assigned transaction identifier.
Instance Method Summary collapse
-
#ack_all ⇒ void
Acknowledges all messages in this poll transaction.
-
#ack_range(sequence_range:) ⇒ void
Acknowledges a specific range of messages by sequence number.
-
#error? ⇒ Boolean
Returns whether the poll resulted in an error.
-
#initialize(transaction_id:, messages:, error:, active_offsets:, transaction_complete:, action_proc: nil) ⇒ QueuePollResponse
constructor
private
A new instance of QueuePollResponse.
-
#nack_all ⇒ void
Negative-acknowledges all messages in this poll transaction.
-
#nack_range(sequence_range:) ⇒ void
Negative-acknowledges a specific range of messages by sequence number.
-
#requeue_all(channel:) ⇒ void
Requeues all messages in this poll transaction to a different channel.
-
#requeue_range(channel:, sequence_range:) ⇒ void
Requeues a specific range of messages to a different channel.
-
#transaction_status ⇒ Boolean
Queries whether the poll transaction is still active on the broker.
Constructor Details
#initialize(transaction_id:, messages:, error:, active_offsets:, transaction_complete:, action_proc: nil) ⇒ QueuePollResponse
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
Returns a new instance of QueuePollResponse.
41 42 43 44 45 46 47 48 49 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 41 def initialize(transaction_id:, messages:, error:, active_offsets:, transaction_complete:, action_proc: nil) @transaction_id = transaction_id @messages = @error = error @active_offsets = active_offsets @transaction_complete = transaction_complete @action_proc = action_proc end |
Instance Attribute Details
#active_offsets ⇒ Array<Integer> (readonly)
Returns currently active message offsets.
32 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 32 attr_reader :transaction_id, :messages, :error, :active_offsets, :transaction_complete |
#error ⇒ String? (readonly)
Returns error description if the poll failed.
32 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 32 attr_reader :transaction_id, :messages, :error, :active_offsets, :transaction_complete |
#messages ⇒ Array<QueueMessageReceived> (readonly)
Returns received messages.
32 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 32 attr_reader :transaction_id, :messages, :error, :active_offsets, :transaction_complete |
#transaction_complete ⇒ Object (readonly)
Returns the value of attribute transaction_complete.
32 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 32 attr_reader :transaction_id, :messages, :error, :active_offsets, :transaction_complete |
#transaction_id ⇒ String (readonly)
Returns broker-assigned transaction identifier.
32 33 34 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 32 def transaction_id @transaction_id end |
Instance Method Details
#ack_all ⇒ void
This method returns an undefined value.
Acknowledges all messages in this poll transaction.
62 63 64 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 62 def ack_all send_action(:AckAll) end |
#ack_range(sequence_range:) ⇒ void
This method returns an undefined value.
Acknowledges a specific range of messages by sequence number.
71 72 73 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 71 def ack_range(sequence_range:) send_action(:AckRange, sequence_range: sequence_range) end |
#error? ⇒ Boolean
Returns whether the poll resulted in an error.
54 55 56 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 54 def error? !@error.nil? && !@error.empty? end |
#nack_all ⇒ void
This method returns an undefined value.
Negative-acknowledges all messages in this poll transaction.
79 80 81 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 79 def nack_all send_action(:NAckAll) end |
#nack_range(sequence_range:) ⇒ void
This method returns an undefined value.
Negative-acknowledges a specific range of messages by sequence number.
88 89 90 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 88 def nack_range(sequence_range:) send_action(:NAckRange, sequence_range: sequence_range) end |
#requeue_all(channel:) ⇒ void
This method returns an undefined value.
Requeues all messages in this poll transaction to a different channel.
97 98 99 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 97 def requeue_all(channel:) send_action(:ReQueueAll, channel: channel) end |
#requeue_range(channel:, sequence_range:) ⇒ void
This method returns an undefined value.
Requeues a specific range of messages to a different channel.
107 108 109 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 107 def requeue_range(channel:, sequence_range:) send_action(:ReQueueRange, channel: channel, sequence_range: sequence_range) end |
#transaction_status ⇒ Boolean
Queries whether the poll transaction is still active on the broker.
rubocop:disable Naming/PredicateMethod -- mirrors broker transaction state, not a Ruby predicate name
115 116 117 118 119 120 121 122 123 |
# File 'lib/kubemq/queues/queue_poll_response.rb', line 115 def transaction_status return false if @transaction_complete return false unless @action_proc response = @action_proc.call(:TransactionStatus) return false if response.nil? !response.TransactionComplete end |