Packages
event_bus
1.0.0-beta2
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
Current section
Files
Jump to
Current section
Files
lib/event_bus/services/watcher.ex
defmodule EventBus.Service.Watcher do
@moduledoc false
alias EventBus.Store
alias :ets, as: Ets
@prefix "eb_ew_"
@doc false
@spec register_topic(String.t | atom()) :: no_return()
def register_topic(topic) do
table_name = table_name(topic)
all_tables = :ets.all()
unless Enum.any?(all_tables, fn table -> table == table_name end) do
opts = [:set, :public, :named_table, {:write_concurrency, true},
{:read_concurrency, true}]
Ets.new(table_name, opts)
end
end
@doc false
@spec unregister_topic(String.t | atom()) :: no_return()
def unregister_topic(topic) do
table_name = table_name(topic)
all_tables = :ets.all()
if Enum.any?(all_tables, fn table -> table == table_name end) do
Ets.delete(table_name)
end
end
@doc false
@spec mark_as_completed(tuple()) :: no_return()
def mark_as_completed({listener, topic, id}) do
{listeners, completers, skippers} = fetch({topic, id})
save_or_delete({topic, id}, {listeners, [listener | completers],
skippers})
end
@doc false
@spec mark_as_skipped(tuple()) :: no_return()
def mark_as_skipped({listener, topic, id}) do
{listeners, completers, skippers} = fetch({topic, id})
save_or_delete({topic, id}, {listeners, completers,
[listener | skippers]})
end
@doc false
@spec fetch(tuple()) :: any()
def fetch({topic, id}) do
case Ets.lookup(table_name(topic), id) do
[{_, data}] -> data
_ -> nil
end
end
@doc false
@spec save(tuple(), tuple()) :: :ok
def save({topic, id}, watcher) do
save_or_delete({topic, id}, watcher)
:ok
end
@spec complete?(tuple()) :: boolean()
defp complete?({listeners, completers, skippers}),
do: length(listeners) == length(completers) + length(skippers)
@spec save_or_delete(tuple(), tuple()) :: no_return()
defp save_or_delete({topic, id}, watcher) do
if complete?(watcher) do
delete_with_relations({topic, id})
else
Ets.insert(table_name(topic), {id, watcher})
end
end
@spec delete_with_relations(tuple()) :: no_return()
defp delete_with_relations({topic, id}) do
Store.delete({topic, id})
Ets.delete(table_name(topic), id)
end
@spec table_name(String.t | atom()) :: atom()
defp table_name(name),
do: :"#{@prefix}#{name}"
end