Packages
ex_esdb
0.0.9-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/consumer.ex
defmodule ExESDB.Repl.Consumer do
@moduledoc false
use GenServer
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
IO.puts("\nCONSUMED #{stream_id}:#{event_type}
version #{inspect(version)}
payload: #{inspect(payload, pretty: true)}")
{:noreply, state}
end
@impl true
def handle_info(_, state) do
{:noreply, state}
end
############## PLUMBING ##############
@impl true
def init(args) do
topic = Keyword.get(args, :topic, "reg_gh:$all")
pubsub = Keyword.get(args, :pubsub, :ex_esdb_pubsub)
pubsub
|> Phoenix.PubSub.subscribe(topic)
{:ok, args}
end
def start_link(args) do
topic = Keyword.get(args, :topic, "reg_gh:$all")
GenServer.start_link(
__MODULE__,
args,
name: Module.concat(__MODULE__, topic)
)
end
def start_consumer(args) do
topic = Keyword.get(args, :topic, "reg_gh:$all")
case start_link(args) do
{:ok, pid} ->
IO.puts("#{Themes.consumer(pid)} for [#{inspect(topic)}] is UP!")
{:error, {:already_started, pid}} ->
IO.puts("#{Themes.consumer(pid)} for [#{inspect(topic)}] is UP!")
{:error, reason} ->
raise "Failed to start consumer for [#{inspect(topic)}].
Reason: #{inspect(reason)}"
end
end
end