Packages
ex_esdb
0.4.0
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/streams_procs.erl
-module(streams_procs).
-include_lib("khepri/include/khepri.hrl").
-export([put_on_new_event/2, get_on_new_event/2]).
on_new_event(Topic) ->
[procs, on_new_event, Topic].
-spec get_on_new_event(Store :: khepri:store(), Id :: string()) -> ok | {error, term()}.
get_on_new_event(Store, Id) when is_atom(Store) ->
Topic = emitter_group:topic(Store, Id),
khepri:get(Store, on_new_event(Topic)).
-spec put_on_new_event(Store :: khepri:store(), Id :: string()) -> ok | {error, term()}.
put_on_new_event(Store, Id) when is_atom(Store) ->
Topic = emitter_group:topic(Store, Id),
case khepri:exists(Store, on_new_event(Topic)) of
true ->
ok;
false ->
ok =
khepri:put(Store,
on_new_event(Topic),
fun(Props) ->
case maps:get(path, Props, undefined) of
undefined -> ok;
Path ->
case streams_store:get_event(Store, Path) of
{ok, undefined} -> ok;
{ok, Event} ->
emitter_group:broadcast(Store, Id, Event),
ok;
{error, Reason} ->
io:format("Broadcasting failed for path ~p to ~p ~n Reason: ~p~n",
[Path, Topic, Reason]),
ok
end
end
end),
% Request asynchronous persistence
ok
end.