Packages
ex_esdb
0.0.10-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
lib/ex_esdb/subscriptions_tracker.ex
defmodule ExESDB.SubscriptionsTracker do
@moduledoc """
As part of the ExESDB.System, the SubscriptionsTracker is responsible for
observing the subscriptions that are maintained in the Store.
Since Khepri triggers are executed on the leader node, the SubscriptionsTracker
will be instructed to start the Emitters system on the leader node whenever a new subscription
is registered.
When a Subscription is deleted, the SubscriptionsTracker will instruct the Emitters system to stop
the associated EmitterPool.
"""
use GenServer
alias ExESDB.Cluster, as: Cluster
alias ExESDB.Emitters, as: Emitters
alias ExESDB.Themes, as: Themes
########### HANDLE_INFO ###########
@impl GenServer
def handle_info({:feature_created, :subscriptions, data}, state) do
IO.puts("Subscription #{inspect(data)} registered")
store = state[:store_id]
if Cluster.leader?(store) do
store
|> Emitters.start_emitter(data)
end
{:noreply, state}
end
@impl GenServer
def handle_info({:feature_deleted, :subscriptions, data}, state) do
IO.puts("Subscription #{inspect(data)} deleted")
# TODO: Stop Emitters
{:noreply, state}
end
def handle_info({:EXIT, pid, reason}, state) do
IO.puts("#{Themes.subscriptions_tracker(pid)} exited with reason: #{inspect(reason)}")
store = state[:store_id]
store
|> :tracker_group.leave(:subscriptions, self())
{:noreply, state}
end
@impl GenServer
def handle_info(_, state) do
{:noreply, state}
end
############## PLUMBING ##############
@impl GenServer
def init(opts) do
Process.flag(:trap_exit, true)
store = Keyword.get(opts, :store_id)
IO.puts("#{Themes.subscriptions_tracker(self())} is UP.")
:ok =
store
|> :subscriptions.setup_tracking(self())
{:ok, opts}
end
@impl true
def terminate(reason, state) do
IO.puts("#{Themes.subscriptions_tracker(self())} terminating with reason: #{inspect(reason)}")
store = state[:store_id]
store
|> :tracker_group.leave(:subscriptions, self())
:ok
end
def child_spec(opts),
do: %{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
type: :worker,
restart: :permanent,
shutdown: 5000
}
def start_link(opts),
do:
GenServer.start_link(
__MODULE__,
opts,
name: __MODULE__
)
end