Packages
ex_esdb
0.0.8-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/monitor.ex
defmodule ExESDB.Repl.EventStreamMonitor do
@moduledoc false
use GenServer
require Logger
alias ExESDB.Themes, as: Themes
alias Phoenix.PubSub, as: PubSub
def start do
opts = ExESDB.Options.app_env()
case start_link(opts) do
{:ok, pid} ->
Logger.info("Monitor started with pid #{inspect(pid)}")
pid
{:error, {:already_started, pid}} ->
IO.puts "Monitor already started with pid #{inspect(pid)}"
pid
{:error, reason} -> raise "Failed to start monitor. Reason: #{inspect(reason)}"
end
end
defp subscribe(opts) do
store_id = inspect(opts[:store_id])
case opts[:pub_sub]
|> PubSub.subscribe(store_id) do
:ok -> Logger.info("#{Themes.monitor(self())} => Subscribed to #{inspect(opts[:store_id], pretty: true)}")
error -> Logger.error("#{Themes.monitor(self())} => Failed to subscribe to #{inspect(opts[:store_id], pretty: true)}. Reason: #{inspect(error)}")
end
end
@impl true
def handle_info({:event_seen, event}, state) do
IO.puts "Seen event #{inspect event}"
{:noreply, state}
end
@impl true
def handle_info(unknown, state) do
IO.puts "Unknown message #{inspect unknown}"
{:noreply, state}
end
@impl true
def init(opts) do
Logger.info("#{Themes.monitor(self())} => Starting monitor for #{inspect(opts[:store_id], pretty: true)}")
subscribe(opts)
{:ok, opts}
end
def start_link(args) do
GenServer.start_link(
__MODULE__,
args,
name: __MODULE__
)
end
def child_spec(opts) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
restart: :permanent,
shutdown: 5000,
type: :worker,
}
end
end