Packages
event_bus
1.1.2
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/subscription.ex
defmodule EventBus.Service.Subscription do
@moduledoc false
alias EventBus.Topic
alias EventBus.Util.Regex, as: RegexUtil
@app :event_bus
@namespace :subscriptions
@doc false
@spec subscribe(tuple()) :: no_return()
def subscribe({listener, topics}) do
{listeners, topic_map} = load_state()
listeners = add_or_update_listener(listeners, {listener, topics})
topic_map =
topic_map
|> add_listener_to_topic_map({listener, topics})
|> Enum.into(%{})
save_state({listeners, topic_map})
end
@doc false
@spec unsubscribe(tuple() | module()) :: no_return()
def unsubscribe(listener) do
{listeners, topic_map} = load_state()
listeners = List.keydelete(listeners, listener, 0)
topic_map =
topic_map
|> remove_listener_from_topic_map(listener)
|> Enum.into(%{})
save_state({listeners, topic_map})
end
@doc false
@spec register_topic(String.t() | atom()) :: no_return()
def register_topic(topic) do
{listeners, topic_map} = load_state()
topic_subscribers =
Enum.reduce(listeners, [], fn {listener, topics}, acc ->
if RegexUtil.superset?(topics, topic), do: [listener | acc], else: acc
end)
save_state({listeners, Map.put(topic_map, topic, topic_subscribers)})
end
@doc false
@spec unregister_topic(String.t() | atom()) :: no_return()
def unregister_topic(topic) do
{listeners, topic_map} = load_state()
save_state({listeners, Map.drop(topic_map, [topic])})
end
@doc false
def subscribers do
{listeners, _topic_map} = load_state()
listeners
end
def subscribers(topic) do
{_listeners, topic_map} = load_state()
topic_map[topic] || []
end
defp remove_listener_from_topic_map(topic_map, listener) do
Enum.map(topic_map, fn {topic, topic_listeners} ->
topic_listeners = List.delete(topic_listeners, listener)
{topic, topic_listeners}
end)
end
defp add_listener_to_topic_map(topic_map, {listener, topics}) do
Enum.map(topic_map, fn {topic, topic_listeners} ->
topic_listeners = List.delete(topic_listeners, listener)
if RegexUtil.superset?(topics, topic) do
{topic, [listener | topic_listeners]}
else
{topic, topic_listeners}
end
end)
end
defp add_or_update_listener(listeners, {listener, topics}) do
if List.keymember?(listeners, listener, 0) do
List.keyreplace(listeners, listener, 0, {listener, topics})
else
[{listener, topics} | listeners]
end
end
defp save_state(state),
do: Application.put_env(@app, @namespace, state, persistent: true)
defp load_state,
do: Application.get_env(@app, @namespace, {[], init_topic_map()})
defp init_topic_map do
topics = Topic.all()
topics
|> Enum.map(fn topic -> {topic, []} end)
|> Enum.into(%{})
end
end