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/ex_esdb/streams_reader.ex
defmodule ExESDB.StreamsReader do
@moduledoc """
This module is responsible for reading events from a stream.
"""
########### API ############
def worker_id(store, stream_id),
do: {:streams_reader_worker, store, stream_id}
@doc """
Returns a list of all streams in the store.
## Parameters
# - `store` is the name of the store.
## Returns
# - `{:ok, streams}` if successful.
"""
@spec get_streams(store :: atom()) :: {:ok, list()} | {:error, term()}
def get_streams(store),
do:
GenServer.call(
get_reader(store, "general_info"),
{:get_streams, store}
)
@doc """
Streams events from `stream` in batches of `count` events, in a `direction`.
"""
@spec stream_events(
store :: atom(),
stream_id :: any(),
start_version :: integer(),
count :: integer(),
direction :: :forward | :backward
) :: {:ok, Enumerable.t()} | {:error, term()}
def stream_events(store, stream_id, start_version, count, direction \\ :forward),
do:
GenServer.call(
get_reader(store, stream_id),
{:stream_events, store, stream_id, start_version, count, direction}
)
defp get_reader(store, stream_id) do
case get_cluster_reader(store, stream_id) do
nil ->
start_reader(store, stream_id)
reader_pid ->
reader_pid
end
end
defp get_cluster_reader(store, stream_id) do
case Swarm.registered()
|> Enum.filter(fn {name, _} ->
match?({:streams_reader_worker, ^store, ^stream_id}, name)
end)
|> Enum.map(fn {_, pid} -> pid end) do
[] ->
nil
readers ->
readers
|> Enum.random()
end
end
defp partition_for(store, stream_id) do
partitions = System.schedulers_online()
key = :erlang.phash2({store, stream_id}, partitions)
key
end
defp start_reader(store, stream_id) do
partition = partition_for(store, stream_id)
case DynamicSupervisor.start_child(
{:via, PartitionSupervisor, {ExESDB.StreamsReaders, partition}},
{ExESDB.StreamsReaderWorker, {store, stream_id, partition}}
) do
{:ok, pid} -> pid
{:error, {:already_started, pid}} -> pid
{:error, reason} -> raise "failed to start streams reader: #{inspect(reason)}"
end
end
end