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
lib/ex_esdb/subscriptions_writer.ex
defmodule ExESDB.SubscriptionsWriter do
@moduledoc """
Provides functions for working with event store subscriptions.
"""
use GenServer
require Logger
alias ExESDB.PersistenceWorker, as: PersistenceWorker
alias ExESDB.Themes, as: Themes
alias ExESDB.StoreNaming
def put_subscription(
store,
type,
selector,
subscription_name \\ "transient",
start_from \\ 0,
subscriber \\ nil
) do
name = StoreNaming.genserver_name(__MODULE__, store)
GenServer.cast(
name,
{:put_subscription, store, type, selector, subscription_name, start_from, subscriber}
)
end
# def put_subscription_sync(
# store,
# type,
# selector,
# subscription_name \\ "transient",
# start_from \\ 0,
# subscriber \\ nil
# ) do
# name = StoreNaming.genserver_name(__MODULE__, store)
# GenServer.call(
# name,
# {:put_subscription_sync, store, type, selector, subscription_name, start_from, subscriber},
# 10_000
# )
# end
#
def delete_subscription(store, type, selector, subscription_name) do
name = StoreNaming.genserver_name(__MODULE__, store)
GenServer.cast(
name,
{:delete_subscription, store, type, selector, subscription_name}
)
end
############ CALLBACKS ############
@impl true
def handle_cast({:delete_subscription, store, type, selector, subscription_name}, state) do
key =
:subscriptions_store.key({type, selector, subscription_name})
if store
|> :khepri.exists!([:subscriptions, key]) do
store
|> :khepri.delete!([:subscriptions, key])
# Request asynchronous persistence
spawn(fn -> PersistenceWorker.request_persistence(store) end)
end
{:noreply, state}
end
@impl GenServer
def handle_cast(
{:put_subscription, store, type, selector, subscription_name, start_from, subscriber},
state
) do
subscription =
%{
selector: selector,
type: type,
subscription_name: subscription_name,
start_from: start_from,
subscriber: subscriber
}
if :subscriptions_store.exists(store, subscription) do
store
|> :subscriptions_store.update_subscription(subscription)
else
store
|> :subscriptions_store.put_subscription(subscription)
end
# Request asynchronous persistence
spawn(fn -> PersistenceWorker.request_persistence(store) end)
{:noreply, state}
end
# @impl GenServer
# def handle_call(
# {:put_subscription_sync, store, type, selector, subscription_name, start_from, subscriber},
# _from,
# state
# ) do
# try do
# subscription =
# %{
# selector: selector,
# type: type,
# subscription_name: subscription_name,
# start_from: start_from,
# subscriber: subscriber
# }
#
# result = if :subscriptions_store.exists(store, subscription) do
# store
# |> :subscriptions_store.update_subscription(subscription)
# else
# store
# |> :subscriptions_store.put_subscription(subscription)
# end
#
# # Request asynchronous persistence
# ExESDB.PersistenceWorker.request_persistence(store)
#
# {:reply, {:ok, result}, state}
# rescue
# error ->
# Logger.warning("Failed to put subscription: #{inspect(error)}")
# {:reply, {:error, error}, state}
# catch
# :exit, reason ->
# Logger.warning("Failed to put subscription (exit): #{inspect(reason)}")
# {:reply, {:error, reason}, state}
# end
# end
#
######## PLUMBING ############
@impl true
def init(opts) do
Process.flag(:trap_exit, true)
IO.puts("#{Themes.subscriptions_writer(self(), "is UP.")}")
{:ok, opts}
end
@impl true
def terminate(reason, _state) do
IO.puts(
"#{Themes.subscriptions_writer(self(), "⚠️ Shutting down gracefully. Reason: #{inspect(reason)}")}"
)
:ok
end
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
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