Packages
ex_esdb
0.5.1
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/subscriptions_procs.erl
-module(subscriptions_procs).
-include_lib("khepri/include/khepri.hrl").
-export([put_on_create_func/1, put_on_delete_func/1, on_create_key/0, on_delete_key/0,
put_on_update_func/1, on_update_key/0]).
on_create_key() ->
[procs, subscriptions, on_create].
on_delete_key() ->
[procs, subscriptions, on_delete].
on_update_key() ->
[procs, subscriptions, on_update].
put_on_create_func(Store) ->
case khepri:exists(Store, on_create_key()) of
true ->
ok;
false ->
ok =
khepri:put(Store,
on_create_key(),
fun(Props) ->
case maps:get(path, Props, undefined) of
undefined -> ok;
Path ->
case subscriptions_store:get_subscription(Store, Path) of
nil -> ok;
Subscription ->
tracker_group:notify_created(Store, subscriptions, Subscription),
ok
end
end
end),
% Request asynchronous persistence
ok;
{error, Reason} ->
io:format("Warning: khepri:exists failed in put_on_create_func with reason: "
"~p~n",
[Reason]),
ok
end.
put_on_update_func(Store) ->
case khepri:exists(Store, on_update_key()) of
true ->
ok;
false ->
ok =
khepri:put(Store,
on_update_key(),
fun(Props) ->
case maps:get(path, Props, undefined) of
undefined -> ok;
Path ->
case subscriptions_store:get_subscription(Store, Path) of
nil -> ok;
Subscription ->
tracker_group:notify_updated(Store, subscriptions, Subscription),
ok
end
end
end),
% Request asynchronous persistence
ok;
{error, Reason} ->
io:format("Warning: khepri:exists failed in put_on_update_func with reason: "
"~p~n",
[Reason]),
ok
end.
put_on_delete_func(Store) ->
case khepri:exists(Store, on_delete_key()) of
true ->
ok;
false ->
ok =
khepri:put(Store,
on_delete_key(),
fun(Props) ->
case maps:get(path, Props, undefined) of
undefined -> ok;
Path ->
case subscriptions_store:get_subscription(Store, Path) of
nil -> ok;
Subscription ->
tracker_group:notify_deleted(Store, subscriptions, Subscription),
ok
end
end
end),
% Request asynchronous persistence
ok;
{error, Reason} ->
io:format("Warning: khepri:exists failed in put_on_delete_func with reason: "
"~p~n",
[Reason]),
ok
end.