Class: Phobos::Actions::ProcessMessage

Inherits:
Object
  • Object
show all
Includes:
Instrumentation
Defined in:
lib/phobos/actions/process_message.rb

Constant Summary

Constants included from Instrumentation

Instrumentation::NAMESPACE

Instance Attribute Summary collapse

Instance Method Summary collapse

Methods included from Instrumentation

#instrument, subscribe, unsubscribe

Constructor Details

#initialize(listener:, message:, listener_metadata:) ⇒ ProcessMessage

Returns a new instance of ProcessMessage.



8
9
10
11
12
13
14
15
16
17
# File 'lib/phobos/actions/process_message.rb', line 8

def initialize(listener:, message:, listener_metadata:)
  @listener = listener
  @message = message
  @metadata = .merge(
    key: message.key,
    partition: message.partition,
    offset: message.offset,
    retry_count: 0
  )
end

Instance Attribute Details

#metadataObject (readonly)

Returns the value of attribute metadata.



6
7
8
# File 'lib/phobos/actions/process_message.rb', line 6

def 
  @metadata
end

Instance Method Details

#executeObject



19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
# File 'lib/phobos/actions/process_message.rb', line 19

def execute
  backoff = @listener.create_exponential_backoff
  payload = force_encoding(@message.value)

  begin
    process_message(payload)
  rescue => e
    retry_count = @metadata[:retry_count]
    interval = backoff.interval_at(retry_count).round(2)

    error = {
      waiting_time: interval,
      exception_class: e.class.name,
      exception_message: e.message,
      backtrace: e.backtrace
    }

    instrument('listener.retry_handler_error', error.merge(@metadata)) do
      Phobos.logger.error do
        { message: "error processing message, waiting #{interval}s" }.merge(error).merge(@metadata)
      end

      sleep interval
    end

    raise Phobos::AbortError if @listener.should_stop?

    @metadata.merge!(retry_count: retry_count + 1)
    retry
  end
end