Packages
ex_esdb
0.4.6
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/streams_reader_worker.ex
defmodule ExESDB.StreamsReaderWorker do
@moduledoc """
Provides functions for reading and streaming events from the event store.
"""
use GenServer
require Logger
alias ExESDB.Themes, as: Themes
alias ExESDB.StreamsHelper, as: Helper
alias ExESDB.StreamsReader, as: StreamsReader
import ExESDB.Khepri.Conditions
defp try_register_with_swarm(name, store, stream_id, partition) do
# Register with Swarm - handle potential registration failures
case register_with_swarm(name) do
:ok ->
{:ok, build_initial_state(name, store, stream_id, partition)}
{:error, reason} ->
Logger.error("Failed to register worker #{inspect(name)}: #{inspect(reason)}")
{:stop, {:registration_failed, reason}}
end
end
defp fetch_single_event(store, stream_id, version) do
# Use 0-based version directly as storage key
padded_version = Helper.pad_version(version, 6)
case :khepri.get(store, [:streams, stream_id, padded_version]) do
{:ok, event} ->
event
{:error, :not_found} ->
nil
{:error, reason} ->
Logger.warning(
"Failed to read event #{version} from stream #{stream_id}: #{inspect(reason)}"
)
nil
end
end
defp stream_events(store, stream_id, start_version, count, direction) do
with :ok <- validate_parameters(start_version, count),
{:ok, event_stream} <- fetch_events(store, stream_id, start_version, count, direction) do
{:ok, event_stream}
else
# If stream doesn't exist, check if that's the real issue
{:error, :stream_not_found} = error ->
case validate_stream_exists(store, stream_id) do
# Stream exists but fetch failed - propagate original error
:ok ->
error
{:error, :stream_not_found} ->
Logger.info(
"Stream #{stream_id} not found, returning empty stream for aggregate loading"
)
# Return empty stream for aggregate loading
{:ok, []}
end
error ->
error
end
rescue
error ->
Logger.error("Error reading stream #{stream_id}: #{inspect(error)}")
{:error, :internal_error}
end
defp validate_stream_exists(store, stream_id) do
exists = Helper.stream_exists?(store, stream_id)
if exists do
Logger.debug("Stream validation: #{stream_id} exists")
:ok
else
Logger.warning("Stream validation: #{stream_id} not found in store #{store}")
# Try to list available streams for debugging
case get_streams_safe(store) do
{:ok, streams} when is_list(streams) and length(streams) > 0 ->
sample_streams = Enum.take(streams, 5)
Logger.warning("Available streams (sample): #{inspect(sample_streams)}")
{:ok, []} ->
Logger.warning("No streams found in store #{store}")
{:error, reason} ->
Logger.warning("Could not list streams: #{inspect(reason)}")
end
{:error, :stream_not_found}
end
end
defp validate_parameters(start_version, count) do
cond do
not is_integer(start_version) -> {:error, {:invalid_start_version, start_version}}
start_version < 0 -> {:error, :invalid_start_version}
not is_integer(count) -> {:error, {:invalid_count, count}}
count < 1 -> {:error, :invalid_count}
true -> :ok
end
end
defp fetch_events(store, stream_id, start_version, count, direction) do
stream_length = Helper.get_version!(store, stream_id)
# Validate stream_length is a proper integer
case stream_length do
length when is_integer(length) and length >= -1 ->
desired_versions = Helper.calculate_versions(start_version, count, direction)
valid_versions =
Enum.filter(desired_versions, fn version ->
version >= 0 && version <= length
end)
event_stream =
valid_versions
|> Stream.map(&fetch_single_event(store, stream_id, &1))
|> Stream.reject(&is_nil/1)
{:ok, event_stream}
invalid_length ->
Logger.warning(
"Invalid stream length #{inspect(invalid_length)} for stream #{stream_id}, treating as empty stream"
)
# Return empty stream instead of error to avoid breaking aggregate loading
{:ok, []}
end
end
############ CALLBACKS ############
@impl true
def handle_call(
{:stream_events, store, stream_id, start_version, count, direction},
_from,
state
) do
result = stream_events(store, stream_id, start_version, count, direction)
{:reply, result, state}
end
@impl true
def handle_call({:get_streams, store}, _from, state) do
result = get_streams_safe(store)
{:reply, result, state}
end
# Handle unexpected call messages gracefully
@impl true
def handle_call(unknown_call, _from, state) do
Logger.warning("Unknown call received: #{inspect(unknown_call)}")
{:reply, {:error, :unknown_call}, state}
end
defp get_streams_safe(store) do
try do
case :khepri.get_many(store, [
:streams,
if_node_exists(exists: true)
]) do
{:ok, stream_data} ->
streams =
Enum.reduce(stream_data, [], fn
{[:streams, stream_id], _stream}, acc
when is_binary(stream_id) or is_atom(stream_id) ->
[stream_id | acc]
{invalid_key, _}, acc ->
Logger.warning("Invalid stream key format: #{inspect(invalid_key)}")
acc
end)
{:ok, Enum.reverse(streams)}
{:error, reason} ->
Logger.error("Failed to get streams: #{inspect(reason)}")
{:error, :store_access_failed}
end
rescue
error ->
Logger.error("Error getting streams: #{inspect(error)}")
{:error, :internal_error}
end
end
################## PlUMBING ##################
@impl true
def init({store, stream_id, partition}) do
Logger.info("[STREAMS_READER_WORKER] Initializing StreamsReaderWorker for store: #{inspect(store)}, stream_id: #{inspect(stream_id)}, partition: #{partition}")
Logger.info("[STREAMS_READER_WORKER] Worker PID: #{inspect(self())}, Node: #{inspect(node())}")
Process.flag(:trap_exit, true)
Logger.info("[STREAMS_READER_WORKER] Process trap_exit enabled for graceful shutdown")
name = StreamsReader.worker_id(store, stream_id)
Logger.info("[STREAMS_READER_WORKER] Generated worker name: #{inspect(name)}")
Logger.info("[STREAMS_READER_WORKER] Attempting to register with Swarm...")
try_register_with_swarm(name, store, stream_id, partition)
try do
# Validate initialization parameters
case validate_init_params(store, stream_id, partition) do
:ok ->
Process.flag(:trap_exit, true)
name = StreamsReader.worker_id(store, stream_id)
# Safe logging - avoid potential crashes from theme formatting
safe_log_startup(name, partition)
try_register_with_swarm(name, store, stream_id, partition)
{:error, reason} ->
Logger.error("Invalid initialization parameters: #{inspect(reason)}")
{:stop, {:invalid_params, reason}}
end
rescue
error ->
Logger.error("Error during worker initialization: #{inspect(error)}")
{:stop, {:init_error, error}}
end
end
defp validate_init_params(store, stream_id, partition) do
cond do
is_nil(store) -> {:error, :store_required}
is_nil(stream_id) or stream_id == "" -> {:error, :stream_id_required}
is_nil(partition) -> {:error, :partition_required}
not is_integer(partition) -> {:error, :partition_must_be_integer}
true -> :ok
end
end
defp safe_log_startup(name, partition) do
try do
msg = "[#{inspect(name)}] is UP on partition #{inspect(partition)}, joining the cluster."
IO.puts("#{Themes.streams_reader_worker(self(), msg)}")
rescue
_error ->
# Fallback to basic logging if theme formatting fails
Logger.info(
"StreamsReaderWorker #{inspect(name)} starting on partition #{inspect(partition)}"
)
end
end
defp register_with_swarm(name) do
try do
case Swarm.register_name(name, self()) do
:yes -> :ok
:no -> {:error, :name_already_registered}
{:error, reason} -> {:error, reason}
other -> {:error, {:unexpected_response, other}}
end
rescue
error -> {:error, {:swarm_error, error}}
end
end
defp build_initial_state(name, store, stream_id, partition) do
%{
worker_name: name,
store: store,
stream_id: stream_id,
node: node(),
partition: partition,
started_at: System.monotonic_time(:millisecond)
}
end
def child_spec({store, stream_id, partition} = args) do
# Validate parameters early to fail fast
case validate_child_spec_params(store, stream_id, partition) do
:ok ->
%{
id: {StreamsReader.worker_id(store, stream_id), partition},
start: {__MODULE__, :start_link, [args]},
type: :worker,
restart: :permanent,
shutdown: 5000
}
{:error, reason} ->
raise ArgumentError, "Invalid child_spec parameters: #{inspect(reason)}"
end
end
defp validate_child_spec_params(store, stream_id, partition) do
cond do
is_nil(store) -> {:error, :store_required}
is_nil(stream_id) or stream_id == "" -> {:error, :stream_id_required}
is_nil(partition) -> {:error, :partition_required}
not is_integer(partition) -> {:error, :partition_must_be_integer}
true -> :ok
end
end
def start_link({store, stream_id, partition} = args) do
# Validate parameters before attempting to start
case validate_child_spec_params(store, stream_id, partition) do
:ok ->
worker_id = StreamsReader.worker_id(store, stream_id)
GenServer.start_link(
__MODULE__,
args,
name: {:global, worker_id}
)
{:error, reason} ->
{:error, {:invalid_params, reason}}
end
end
# Handle unexpected info messages
@impl true
def handle_info(msg, state) do
Logger.debug("Unexpected info message: #{inspect(msg)}")
{:noreply, state}
end
# Handle unexpected cast messages
@impl true
def handle_cast(msg, state) do
Logger.warning("Unexpected cast message: #{inspect(msg)}")
{:noreply, state}
end
# Graceful termination
@impl true
def terminate(reason, state) do
worker_name = Map.get(state, :worker_name, "unknown")
case reason do
:normal ->
Logger.info("StreamsReaderWorker #{inspect(worker_name)} terminating normally")
:shutdown ->
Logger.info("StreamsReaderWorker #{inspect(worker_name)} shutting down")
{:shutdown, _} ->
Logger.info("StreamsReaderWorker #{inspect(worker_name)} shutting down")
_ ->
Logger.warning(
"StreamsReaderWorker #{inspect(worker_name)} terminating: #{inspect(reason)}"
)
end
:ok
end
end