Packages
ex_esdb
0.0.19
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/gateway_worker.ex
defmodule ExESDB.GatewayWorker do
@moduledoc """
GatewayWorker processes are started on each node in the cluster,
and contain the implementation functions for the GatewayAPI.
"""
use GenServer
alias ExESDB.SnapshotsReader, as: SnapshotsR
alias ExESDB.SnapshotsWriter, as: SnapshotsW
alias ExESDB.SubscriptionsReader, as: SubsR
alias ExESDB.SubscriptionsWriter, as: SubsW
alias ExESDB.StreamsHelper, as: StreamsH
alias ExESDB.StreamsReader, as: StreamsR
alias ExESDB.StreamsWriter, as: StreamsW
alias ExESDB.Themes, as: Themes
require Logger
@type store :: atom()
@type stream :: String.t()
@type subscription_name :: String.t()
@type error :: term
@type subscription_type :: :by_stream | :by_event_type | :by_event_pattern
@type selector_type :: String.t() | map()
############ HANDLE_CALL ############
@impl GenServer
def handle_call({:stream_forward, store, stream_id, start_version, count}, _from, state) do
case store
|> StreamsR.stream_events(stream_id, start_version, count, :forward) do
{:ok, event_stream} ->
{:reply, {:ok, event_stream}, state}
{:error, reason} ->
{:reply, {:error, reason}, state}
end
end
@impl GenServer
def handle_call({:stream_backward, store, stream_id, start_version, count}, _from, state) do
case store
|> StreamsR.stream_events(stream_id, start_version, count, :backward) do
{:ok, event_stream} ->
{:reply, {:ok, event_stream}, state}
{:error, reason} ->
{:reply, {:error, reason}, state}
end
end
@impl GenServer
def handle_call({:get_events, store, stream_id, start_version, count, direction}, _from, state) do
case store
|> StreamsR.stream_events(stream_id, start_version, count, direction) do
{:ok, events} ->
{:reply, {:ok, events |> Enum.to_list()}, state}
{:error, reason} ->
{:reply, {:error, reason}, state}
end
end
@impl GenServer
def handle_call({:get_streams, store}, _from, state) do
case store
|> StreamsR.get_streams() do
{:ok, streams} ->
{:reply, {:ok, streams |> Enum.to_list()}, state}
{:error, _reason} ->
{:reply, {:ok, []}, state}
end
end
@impl GenServer
def handle_call({:get_subscriptions, store}, _from, state) do
reply =
store
|> SubsR.get_subscriptions()
{:reply, {:ok, reply}, state}
end
@impl GenServer
def handle_call({:get_version, store, stream}, _from, state) do
version =
store
|> StreamsH.get_version!(stream)
{:reply, {:ok, version}, state}
end
@impl GenServer
def handle_call({:append_events, store, stream_id, events}, _from, state) do
current_version =
store
|> StreamsH.get_version!(stream_id)
case store
|> StreamsW.append_events(stream_id, current_version, events) do
{:ok, new_version} ->
{:reply, {:ok, new_version}, state}
{:error, reason} ->
{:reply, {:error, reason}, state}
end
end
@impl GenServer
def handle_call({:read_snapshot, store, source_uuid, stream_uuid, version}, _from, state) do
case store
|> SnapshotsR.read_snapshot(source_uuid, stream_uuid, version) do
{:ok, snapshot} ->
{:reply, {:ok, snapshot}, state}
{:error, reason} ->
{:reply, {:error, reason}, state}
end
end
@impl GenServer
def handle_call({:list_snapshots, store, source_uuid, stream_uuid}, _from, state) do
case store
|> SnapshotsR.list_snapshots(source_uuid, stream_uuid) do
{:ok, snapshots} ->
{:reply, {:ok, snapshots}, state}
{:error, reason} ->
{:reply, {:error, reason}, state}
end
end
################ HANDLE_CAST #############
@impl true
def handle_cast(
{:remove_subscription, store, type, selector, subscription_name},
state
) do
store
|> SubsW.delete_subscription(type, selector, subscription_name)
{:noreply, state}
end
@impl true
def handle_cast(
{:save_subscription, store, type, selector, subscription_name, start_from, subscriber},
state
) do
store
|> SubsW.put_subscription(type, selector, subscription_name, start_from, subscriber)
{:noreply, state}
end
@impl true
def handle_cast(
{:ack_event, store, subscription_name, subscriber_pid, event},
state
) do
%{
event_stream_id: stream_id,
event_number: event_number
} = event
store
|> SubsW.put_subscription(
:by_stream,
"$#{stream_id}",
subscription_name,
event_number + 1,
subscriber_pid
)
{:noreply, state}
end
@impl true
def handle_cast(
{:record_snapshot, store, source_uuid, stream_uuid, version, snapshot_record},
state
) do
store
|> SnapshotsW.record_snapshot(source_uuid, stream_uuid, version, snapshot_record)
{:noreply, state}
end
@impl true
def handle_cast({:delete_snapshot, store, source_uuid, stream_uuid, version}, state) do
store
|> SnapshotsW.delete_snapshot(source_uuid, stream_uuid, version)
{:noreply, state}
end
############# PLUMBING #############
def child_spec(opts) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
type: :worker,
restart: :permanent,
shutdown: 5000
}
end
def start_link(opts) do
GenServer.start_link(
__MODULE__,
opts,
name: __MODULE__
)
end
def gateway_worker_name,
do: {:gateway_worker, node(), :rand.uniform(10_000)}
@impl true
def init(opts) do
Process.flag(:trap_exit, true)
name = gateway_worker_name()
new_state = Keyword.put(opts, :gateway_worker_name, name)
msg = "[#{inspect(name)}] is UP, joining the cluster."
IO.puts("#{Themes.gateway_worker(msg)}")
Swarm.register_name(name, self())
{:ok, new_state}
end
@impl true
def terminate(reason, state) do
name = Keyword.get(state, :gateway_worker_name)
msg = "[#{inspect(name)}] is TERMINATED with reason #{inspect(reason)}, leaving the cluster."
IO.puts("#{Themes.gateway_worker(msg)}")
Swarm.unregister_name(name)
:ok
end
@impl true
def handle_info({:EXIT, _pid, reason}, state) do
name = Keyword.get(state, :gateway_worker_name)
msg = "[#{inspect(name)}] is EXITING with reason #{inspect(reason)}, leaving the cluster."
IO.puts("#{Themes.gateway_worker(msg)}")
Swarm.unregister_name(name)
{:noreply, state}
end
end