Packages
ex_esdb
0.4.2
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
lib/ex_esdb/subscriptions_reader.ex
defmodule ExESDB.SubscriptionsReader do
@moduledoc """
Provides functions for working with event store subscriptions.
"""
use GenServer
import ExESDB.Khepri.Conditions
alias ExESDB.Themes, as: Themes
alias ExESDB.StoreNaming
require Logger
def get_subscriptions(store) do
name = StoreNaming.genserver_name(__MODULE__, store)
GenServer.call(
name,
{:get_subscriptions, store}
)
end
################ HANDLE_CALL #############
@impl GenServer
def handle_call({:get_subscriptions, store}, _from, state) do
case store
|> :khepri.get_many([
:subscriptions,
if_all(
conditions: [
if_path_matches(regex: :any),
if_has_payload(has_payload: true)
]
)
]) do
{:ok, result} ->
{:reply, result, state}
_ ->
{:reply, [], state}
end
end
############### PLUMBING ###############
def start_link(opts) do
store_id = StoreNaming.extract_store_id(opts)
name = StoreNaming.genserver_name(__MODULE__, store_id)
GenServer.start_link(
__MODULE__,
opts,
name: name
)
end
@impl true
def init(opts) do
Process.flag(:trap_exit, true)
IO.puts("#{Themes.subscriptions_reader(self(), "is UP")}")
{:ok, opts}
end
@impl true
def terminate(reason, _state) do
IO.puts("#{Themes.subscriptions_reader(self(), "⚠️ Shutting down gracefully. Reason: #{inspect(reason)}")}")
:ok
end
def child_spec(opts) do
store_id = StoreNaming.extract_store_id(opts)
%{
id: StoreNaming.child_spec_id(__MODULE__, store_id),
start: {__MODULE__, :start_link, [opts]},
type: :worker,
restart: :permanent,
shutdown: 5000
}
end
end