Packages
event_bus
1.3.1
1.7.0
1.6.2
1.6.1
1.6.0
retired
1.5.2
retired
1.5.1
retired
1.5.0
retired
1.4.4
retired
1.4.3
retired
1.4.2
retired
1.4.1
retired
1.4.0
retired
1.3.8
retired
1.3.7
retired
1.3.6
retired
1.3.5
retired
1.3.4
retired
1.3.3
retired
1.3.2
retired
1.3.1
retired
1.3.0
retired
1.2.0
retired
1.1.3
retired
1.1.2
retired
1.1.1
retired
1.1.0
retired
1.0.0
retired
1.0.0-beta4
retired
1.0.0-beta3
retired
1.0.0-beta2
retired
1.0.0-beta1
retired
0.9.0
retired
0.8.0
retired
0.7.0
retired
0.6.1
retired
0.6.0
retired
0.5.0
retired
0.4.1
retired
0.4.0
retired
0.3.1
retired
0.3.0
retired
0.2.1
retired
0.2.0
retired
0.1.0
retired
Traceable, extendable and minimalist event bus implementation for Elixir with built-in event store and event watcher based on ETS
Retired package: Deprecated - Please prefer always latest minor
Current section
Files
Jump to
Current section
Files
lib/event_bus/services/notification.ex
defmodule EventBus.Service.Notification do
@moduledoc false
require Logger
alias EventBus.Manager.{Observation, Store, Subscription}
alias EventBus.Model.Event
@logging_level :info
@doc false
@spec notify(Event.t()) :: no_return()
def notify(%Event{id: id, topic: topic} = event) do
listeners = Subscription.subscribers(topic)
:ok = Store.create(event)
:ok = Observation.create({listeners, topic, id})
notify_listeners(listeners, {topic, id})
end
@spec notify_listeners(list(), tuple()) :: no_return()
defp notify_listeners(listeners, event_shadow) do
for listener <- listeners, do: notify_listener(listener, event_shadow)
end
@spec notify_listener(tuple(), tuple()) :: no_return()
@spec notify_listener(module(), tuple()) :: no_return()
defp notify_listener({listener, config}, {topic, id}) do
listener.process({config, topic, id})
rescue
error ->
log(listener, error)
Observation.mark_as_skipped({{listener, config}, topic, id})
end
defp notify_listener(listener, {topic, id}) do
listener.process({topic, id})
rescue
error ->
log(listener, error)
Observation.mark_as_skipped({listener, topic, id})
end
@spec log(module(), any()) :: no_return()
defp log(listener, error) do
msg = "#{listener}.process/1 raised an error!\n#{inspect(error)}"
Logger.log(@logging_level, msg)
end
end