Module: Julewire::Karafka::MessageContext

Defined in:
lib/julewire/karafka/message_context.rb

Class Method Summary collapse

Class Method Details

.call(message, configuration:, fields: nil) ⇒ Object



7
8
9
10
# File 'lib/julewire/karafka/message_context.rb', line 7

def call(message, configuration:, fields: nil, &)
  fields ||= PayloadReader.message_payload(message)
  call_fields(fields, configuration: configuration, &)
end

.call_fields(fields, configuration:) ⇒ Object



12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
# File 'lib/julewire/karafka/message_context.rb', line 12

def call_fields(fields, configuration:, &)
  carrier = carrier_for(fields, configuration)

  result = Core::Propagation::Carrier.extract_result(
    carrier,
    key: configuration.carrier_key,
    max_bytes: configuration.carrier_max_bytes
  )
  record_carrier_restore_failure(result)
  fields = Core::Fields::FieldSet.deep_symbolize_keys(fields)

  Core::Propagation.restore(result.envelope, owned: true) do
    Core::Integration::Facade.with_neutral(message_neutral(fields)) do
      Core::Integration::Facade.with_attributes(message_attributes(fields), &)
    end
  end
end