Packages
ex_esdb
0.4.8
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_store.erl
-module(subscriptions_store).
-include_lib("khepri/include/khepri.hrl").
-export([get_subscription/2, put_subscription/2, delete_subscription/2, key/1,
update_subscription/2, exists/2]).
get_subscription(Store, Key) when is_binary(Key) ->
Path = [subscriptions, Key],
get_subscription(Store, Path);
get_subscription(Store, Path) when is_list(Path) ->
case khepri:get(Store, Path) of
{ok, Subscription} ->
Subscription;
{error, _Reason} ->
nil
end.
key(KeyT) when is_tuple(KeyT) ->
integer_to_binary(erlang:phash2(KeyT));
key(Subscription) when is_map(Subscription) ->
#{type := Type,
selector := Selector,
subscription_name := SubscriptionName} =
Subscription,
key({Type, Selector, SubscriptionName}).
-spec exists(atom(), map()) -> boolean().
exists(Store, Subscription) ->
Key = key(Subscription),
case khepri:exists(Store, [subscriptions, Key]) of
true ->
true;
false ->
false;
{error, Reason} ->
io:format("Warning: khepri:exists failed in exists/2 with reason: ~p~n", [Reason]),
false
end.
-spec put_subscription(atom(), map()) -> ok.
put_subscription(Store, Subscription) ->
Key = key(Subscription),
case khepri:exists(Store, [subscriptions, Key]) of
true ->
ok = khepri:update(Store, [subscriptions, Key], Subscription),
% Request asynchronous persistence
ok;
false ->
ok = khepri:put(Store, [subscriptions, Key], Subscription),
% Request asynchronous persistence
ok
end.
-spec delete_subscription(atom(), map()) -> ok.
delete_subscription(Store, Subscription) ->
Key = key(Subscription),
case khepri:exists(Store, [subscriptions, Key]) of
true ->
ok = khepri:delete(Store, [subscriptions, Key]),
% Request asynchronous persistence
ok;
false ->
ok
end.
-spec update_subscription(atom(), map()) -> ok.
update_subscription(Store, Subscription) ->
Key = key(Subscription),
case khepri:exists(Store, [subscriptions, Key]) of
true ->
ok = khepri:update(Store, [subscriptions, Key], Subscription),
% Request asynchronous persistence
ok;
false ->
ok = khepri:put(Store, [subscriptions, Key], Subscription),
% Request asynchronous persistence
ok
end.