Packages
ex_esdb
0.0.9-alpha
0.11.0
0.10.0
0.9.0
0.8.0
0.7.8
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.1
0.6.0
0.5.1
0.5.0
0.4.8
0.4.7
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.3
0.3.2
0.3.1
0.3.0
0.2.5
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
0.0.20
0.0.19
0.0.18
0.0.17
0.0.16
0.0.15
0.0.14-alpha
0.0.13-alpha
0.0.12-alpha
0.0.11-alpha
0.0.10-alpha
0.0.9-alpha
0.0.8-alpha
0.0.6-alpha
0.0.5-alpha
0.0.4-alpha
0.0.3-alpha
0.0.2-alfa
0.0.1-alfa
ExESDB is a reincarnation of rabbitmq/khepri, specialized for use as a BEAM-native event store.
Current section
Files
Jump to
Current section
Files
src/emitter_group.erl
-module(emitter_group).
-export([join/3, members/2, broadcast/3, group_key/2, topic/2, emitter_name/2,
emitter_name/3, persist_emitters/3, setup_emitters/4]).
-spec setup_emitters(Store :: khepri:store(),
Id :: string(),
Filter :: khepri:filter(),
PoolSize :: integer()) ->
list().
setup_emitters(Store, Id, Filter, PoolSize) ->
ok = streams:setup_when_new_event(Store, Id, Filter),
Emitters = persist_emitters(Store, Id, PoolSize),
Emitters.
-spec join(Store :: atom(), Id :: string(), PidOrPids :: pid() | [pid()]) -> ok.
join(Store, Id, PidOrPids) when is_atom(Store) ->
Group = group_key(Store, Id),
ok = pg:join('Elixir.Phoenix.PubSub', Group, PidOrPids),
ok.
-spec members(Store :: atom(), Id :: string()) -> [pid()].
members(Store, Id) when is_atom(Store) ->
Group = group_key(Store, Id),
pg:get_members('Elixir.Phoenix.PubSub', Group).
-spec random_emitter(Emitters :: [pid()]) -> {error, {no_such_group, term()}} | pid().
random_emitter(Emitters) ->
EmittersTuple = list_to_tuple(Emitters),
Size = tuple_size(EmittersTuple),
Random = rand:uniform(Size),
erlang:element(Random, EmittersTuple).
-spec broadcast(Store :: atom(), Id :: string(), Event :: map()) -> ok | {error, term()}.
broadcast(Store, Id, Event) when is_atom(Store) ->
Topic = topic(Store, Id),
case random_emitter(members(Store, Id)) of
{error, {no_such_group, _}} ->
logger:error("NO_GROUP [~p]~n", [Topic]),
{error, no_such_group};
EmitterPid ->
Message =
if node(EmitterPid) =:= node() ->
forward_to_local_msg(Topic, Event);
true ->
broadcast_msg(Topic, Event)
end,
EmitterPid ! Message
end.
-spec forward_to_local_msg(Topic :: binary(), Event :: map()) -> tuple().
forward_to_local_msg(Topic, Event) ->
{forward_to_local, Topic, Event}.
-spec broadcast_msg(Topic :: binary(), Event :: map()) -> tuple().
broadcast_msg(Topic, Event) ->
{broadcast, Topic, Event}.
group_key(Store, Id) ->
{Store, Id, emitters}.
-spec topic(Store :: atom(), Id :: string()) -> binary().
topic(Store, <<"$all">>) ->
iolist_to_binary(io_lib:format("~s:$all", [Store]));
topic(Store, Id) ->
iolist_to_binary(io_lib:format("~s:~s", [Store, Id])).
-spec emitter_name(Store :: atom(), Id :: string()) -> atom().
emitter_name(Store, Id) ->
list_to_atom(lists:flatten(
io_lib:format("~s_~s_emitter", [Store, Id]))).
-spec emitter_name(Store :: atom(), Id :: string(), Number :: integer()) -> atom().
emitter_name(Store, Id, Number) ->
list_to_atom(lists:flatten(
io_lib:format("~s_~s_emitter_~p", [Store, Id, Number]))).
-spec persist_emitters(Store :: atom(), Id :: string(), PoolSize :: integer()) -> list().
persist_emitters(Store, Id, PoolSize) ->
% Generate a list of emitter names
EmitterList = [emitter_name(Store, Id, Number) || Number <- lists:seq(1, PoolSize)],
Emitters = [emitter_name(Store, Id) | EmitterList],
Key = group_key(Store, Id),
persistent_term:put(Key, list_to_tuple(Emitters)),
Emitters.
-spec retrieve_emitters(Store :: atom(), Id :: string()) -> list().
retrieve_emitters(Store, Id) ->
Key = group_key(Store, Id),
EmitterTuple = persistent_term:get(Key),
tuple_to_list(EmitterTuple).