Packages
ex_esdb
0.4.0
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/emitter_worker.ex
defmodule ExESDB.EmitterWorker do
@moduledoc """
As part of the ExESDB.System,
the EmitterWorker is responsible for managing the communication
between the Event Store and the PubSub mechanism.
"""
use GenServer
alias ExESDB.Options, as: Options
alias Phoenix.PubSub, as: PubSub
require ExESDB.Themes, as: Themes
require Logger
defp send_or_kill_pool(pid, event, store, selector) do
if Process.alive?(pid) do
Process.send(pid, {:events, [event]}, [])
else
ExESDB.EmitterPool.stop(store, selector)
end
end
defp emit(pub_sub, topic, event) do
pub_sub
|> PubSub.broadcast(topic, {:events, [event]})
end
@impl GenServer
def init({store, sub_topic, subscriber}) do
Logger.info("[EMITTER_WORKER] Initializing EmitterWorker for store: #{inspect(store)}, topic: #{inspect(sub_topic)}")
Logger.info("[EMITTER_WORKER] EmitterWorker PID: #{inspect(self())}, Node: #{inspect(node())}")
Logger.info("[EMITTER_WORKER] Subscriber: #{inspect(subscriber)}")
Process.flag(:trap_exit, true)
Logger.info("[EMITTER_WORKER] Process trap_exit enabled for graceful shutdown")
scheduler_id = :erlang.system_info(:scheduler_id)
topic = :emitter_group.topic(store, sub_topic)
Logger.info("[EMITTER_WORKER] Generated topic: #{inspect(topic)}")
Logger.info("[EMITTER_WORKER] Running on scheduler: #{inspect(scheduler_id)}")
Logger.info("[EMITTER_WORKER] Joining emitter group for store: #{inspect(store)}, topic: #{inspect(sub_topic)}")
:ok = :emitter_group.join(store, sub_topic, self())
Logger.info("[EMITTER_WORKER] ✅ Successfully joined emitter group")
msg = "for #{inspect(topic)} is UP on scheduler #{inspect(scheduler_id)}"
IO.puts("#{Themes.emitter_worker(self(), msg)}")
Logger.info("[EMITTER_WORKER] EmitterWorker initialization complete")
{:ok, %{subscriber: subscriber, store: store, selector: sub_topic}}
end
@impl GenServer
def terminate(reason, %{store: store, selector: selector}) do
msg = "is TERMINATED with reason #{inspect(reason)}"
IO.puts("#{Themes.emitter_worker(self(), msg)}")
:ok = :emitter_group.leave(store, selector, self())
:ok
end
def start_link({store, sub_topic, subscriber, emitter}),
do:
GenServer.start_link(
__MODULE__,
{store, sub_topic, subscriber},
name: emitter
)
def child_spec({store, sub_topic, subscriber, emitter}) do
%{
id: Module.concat(__MODULE__, emitter),
start: {__MODULE__, :start_link, [{store, sub_topic, subscriber, emitter}]},
restart: :permanent,
shutdown: 5000,
type: :worker
}
end
@impl true
def handle_info(
{:broadcast, topic, event},
%{subscriber: subscriber, store: store, selector: selector} = state
) do
case subscriber do
nil ->
pubsub = Options.pub_sub()
pubsub
|> emit(topic, event)
pid ->
send_or_kill_pool(pid, event, store, selector)
end
{:noreply, state}
end
@impl true
def handle_info(
{:forward_to_local, topic, event},
%{subscriber: subscriber, store: store, selector: selector} = state
) do
case subscriber do
nil ->
pubsub = Options.pub_sub()
pubsub
|> emit(topic, event)
pid ->
send_or_kill_pool(pid, event, store, selector)
end
{:noreply, state}
end
@impl true
def handle_info({:events, events}, state) when is_list(events) do
# Handle events messages - these might come from feedback loops or external systems
# Just ignore them since they're already processed events
{:noreply, state}
end
@impl true
def handle_info(msg, state) do
Logger.warning("Received unexpected message #{inspect(msg)} on #{inspect(self())}")
{:noreply, state}
end
@impl GenServer
def handle_cast({:update_subscriber, new_subscriber}, state) do
Logger.info("EmitterWorker updating subscriber from #{inspect(state.subscriber)} to #{inspect(new_subscriber)}")
updated_state = %{state | subscriber: new_subscriber}
{:noreply, updated_state}
end
end