Packages
ex_esdb
0.0.14-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/repl/observer.ex
defmodule ExESDB.Repl.Observer do
@moduledoc """
The Repl.Observer is a GenServer that:
- adds a transient subscription to the store.
- subscribes to the events emitted by the store, via Phoenix PubSub.
- prints the events to the console.
"""
use GenServer
require Logger
alias ExESDB.GatewayAPI, as: API
alias ExESDB.Options, as: Options
alias ExESDB.Themes, as: Themes
@impl true
def handle_info({:event_emitted, event}, state) do
%{
event_stream_id: stream_id,
event_type: event_type,
event_number: version,
data: payload
} = event
msg = "#{stream_id}:#{event_type} (v#{version}) => #{inspect(payload, pretty: true)}"
IO.puts(Themes.observed(msg))
{:noreply, state}
end
@impl true
def handle_info(msg, state) do
Logger.error("Received unexpected message #{inspect(msg)}")
{:noreply, state}
end
############## PLUMBING ##############
@impl true
def init(args) do
store = store(args)
pubsub = pubsub(args)
selector = selector(args)
type = type(args)
name = name(args)
topic = topic(store, selector, name)
:ok =
store
|> API.save_subscription(type, selector, name)
:ok =
pubsub
|> Phoenix.PubSub.subscribe(topic)
{:ok, args}
end
def start_link(args) do
store = store(args)
pubsub = pubsub(args)
selector = selector(args)
name = name(args)
topic = topic(store, selector, name)
args =
args
|> Keyword.put(:store, store)
|> Keyword.put(:pubsub, pubsub)
GenServer.start_link(
__MODULE__,
args,
name: Module.concat(__MODULE__, topic)
)
end
@spec start(keyword()) :: pid()
@doc """
Starts an observer process for a given topic.
## Parameters
- `store`: The store to consume events from (atom, default: the configured store).
- `type`: The type of subscription to consume events from (atom, default: `:by_stream`).
- `selector`: The selector of the subscription to consume events from (string, default: `"$all"`).
- `topic`: The topic to consume events from (string, default: `reg_gh:$all`).
- `name`: The name of the observer (string, default: `transient`).
"""
def start(args) do
store = store(args)
selector = selector(args)
name = name(args)
topic = topic(store, selector, name)
case start_link(args) do
{:ok, pid} ->
IO.puts("#{Themes.observer(pid)} for [#{inspect(topic)}] is UP!")
pid
{:error, {:already_started, pid}} ->
IO.puts("#{Themes.observer(pid)} for [#{inspect(topic)}] is UP!")
pid
{:error, reason} ->
raise "Failed to start observer for [#{inspect(topic)}].
Reason: #{inspect(reason)}"
end
end
defp store(args), do: Keyword.get(args, :store, Options.store_id())
defp type(args), do: Keyword.get(args, :type, :by_stream)
defp selector(args), do: Keyword.get(args, :selector, "$all")
defp pubsub(args), do: Keyword.get(args, :pubsub, Options.pub_sub())
defp name(args), do: Keyword.get(args, :name, "transient")
defp topic(store, selector, "transient"), do: :emitter_group.topic(store, selector)
defp topic(store, _, name), do: :emitter_group.topic(store, name)
end