Packages
kafka_message_bus
4.2.10
4.3.4
4.3.3
4.3.2
4.3.1
4.3.0
4.2.10
4.2.9
4.2.8
4.2.7
4.2.6
4.2.5
4.2.4
4.2.3
4.2.2
4.2.1
4.2.0
4.1.1
4.1.0
4.0.1
4.0.0
4.0.0-rc.8
4.0.0-rc.7
4.0.0-rc.6
4.0.0-rc.5
4.0.0-rc.4
4.0.0-rc.3
4.0.0-rc.2
3.0.0
2.2.4
2.2.3
2.2.2
2.2.1
2.2.0
2.1.0
2.0.2
2.0.1
2.0.0
1.0.1
1.0.0
0.2.3
0.2.2
0.2.1
0.2.0
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
A general purpose messaging utility library supporting Exq and Kaffe
Current section
Files
Jump to
Current section
Files
lib/kafka_message_bus/consumer_handler.ex
defmodule KafkaMessageBus.ConsumerHandler do
@moduledoc """
Contains perform/2 function that is used to consume a message using the
specified module.
"""
alias KafkaMessageBus.{MessageDataValidator, Utils}
require Logger
def perform(module, message) when is_binary(module) do
module = String.to_existing_atom(module)
perform(module, message)
end
def perform(module, message) do
Logger.debug(fn ->
"Running #{Utils.to_module_short_name(module)} message handler"
end)
case MessageDataValidator.validate(message) do
{:ok, :message_contract_excluded} ->
Logger.info(fn ->
"Message contract (consume) excluded: message=#{inspect(message)}"
end)
consume(module, message)
{:ok, message_data} ->
m = Map.put(message, "data", message_data)
consume(module, m)
{:error, [] = validation_errors} ->
Logger.warn(fn ->
"Validation failed for message_data consumption: #{inspect(validation_errors)}\n#{
inspect(message)
}"
end)
{:error, validation_errors}
{:error, :unrecognized_message_data_type} ->
Logger.error(fn ->
"Attempting to consume unrecognized message data type: #{inspect(message)}"
end)
# DEPRECATED: We currently try to consume messages that are not recognized by
# the validator but want an error to be returned in the future. This may be
# changed once all the apps that depend on kafka_message_bus are updated to
# support the currently known message contracts.
consume(module, message)
err ->
Logger.error(fn ->
"Unexpected response encountered when validating consumer message data: #{inspect(err)}"
end)
err
end
rescue
e ->
trace = Exception.format_stacktrace(__STACKTRACE__)
err_msg =
"Failed to process module. module: #{inspect(module)}" <>
", error: #{inspect(e)}" <> ", trace: #{trace}" <> ", message: #{inspect(message)}"
Logger.error(fn -> err_msg end)
{:error, e}
end
defp consume(module, message) do
case module.process(message) do
:ok ->
:ok
{:ok, _} ->
:ok
{:error, reason} = result ->
Logger.error(fn ->
"Error encountered while executing module.process/1: #{inspect(reason)}"
end)
result
other ->
Logger.warn(fn ->
"Unexpected response when executing module.process/1: #{inspect(other)}"
end)
{:error, other}
end
end
end