Class: KubeMQ::Queues::QueuePollResponse

Inherits:
Object
  • Object
show all
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.

Examples:

Acknowledge all messages

response = client.poll(request)
unless response.error?
  response.messages.each { |msg| process(msg.body) }
  response.ack_all
end

See Also:

Instance Attribute Summary collapse

Instance Method Summary collapse

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.

Parameters:

  • transaction_id (String)

    broker-assigned transaction ID

  • messages (Array<QueueMessageReceived>)

    received messages

  • error (String, nil)

    error description on failure

  • active_offsets (Array<Integer>)

    active message offsets

  • transaction_complete (Boolean)

    whether the transaction is complete

  • action_proc (Proc, nil) (defaults to: nil)

    internal callback for transaction actions



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 = messages
  @error = error
  @active_offsets = active_offsets
  @transaction_complete = transaction_complete
  @action_proc = action_proc
end

Instance Attribute Details

#active_offsetsArray<Integer> (readonly)

Returns currently active message offsets.

Returns:

  • (Array<Integer>)

    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

#errorString? (readonly)

Returns error description if the poll failed.

Returns:

  • (String, nil)

    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

#messagesArray<QueueMessageReceived> (readonly)

Returns received messages.

Returns:



32
# File 'lib/kubemq/queues/queue_poll_response.rb', line 32

attr_reader :transaction_id, :messages, :error, :active_offsets, :transaction_complete

#transaction_completeObject (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_idString (readonly)

Returns broker-assigned transaction identifier.

Returns:

  • (String)

    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_allvoid

This method returns an undefined value.

Acknowledges all messages in this poll transaction.

Raises:

  • (TransactionError)

    if no downstream receiver is bound or the broker rejects the action



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.

Parameters:

  • sequence_range (Array<Integer>)

    sequence numbers to acknowledge

Raises:

  • (TransactionError)

    if no downstream receiver is bound or the broker rejects the action



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.

Returns:

  • (Boolean)

    true if an error occurred



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

def error?
  !@error.nil? && !@error.empty?
end

#nack_allvoid

This method returns an undefined value.

Negative-acknowledges all messages in this poll transaction.

Raises:

  • (TransactionError)

    if no downstream receiver is bound or the broker rejects the action



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.

Parameters:

  • sequence_range (Array<Integer>)

    sequence numbers to nack

Raises:

  • (TransactionError)

    if no downstream receiver is bound or the broker rejects the action



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.

Parameters:

  • channel (String)

    destination channel name

Raises:

  • (TransactionError)

    if no downstream receiver is bound or the broker rejects the action



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.

Parameters:

  • channel (String)

    destination channel name

  • sequence_range (Array<Integer>)

    sequence numbers to requeue

Raises:

  • (TransactionError)

    if no downstream receiver is bound or the broker rejects the action



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_statusBoolean

Queries whether the poll transaction is still active on the broker.

rubocop:disable Naming/PredicateMethod -- mirrors broker transaction state, not a Ruby predicate name

Returns:

  • (Boolean)

    true if the transaction is still active, false if completed or no receiver



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