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/streams_writer.ex
defmodule ExESDB.StreamsWriter do
@moduledoc """
This module is responsible for writing events to a stream.
It is actually an API style wrapper around the StreamsWriterWorker.
"""
########### API ############
@spec append_events(
store :: atom(),
stream_id :: any(),
expected_version :: integer(),
events :: list()
) :: {:ok, integer()} | {:error, term()}
def append_events(store, stream_id, expected_version, events) do
writer = get_writer(store, stream_id)
GenServer.call(
writer,
{:append_events, store, stream_id, expected_version, events}
)
end
def worker_id(store, stream_id),
do: {:streams_writer_worker, store, stream_id}
def hr_worker_id_atom(store, stream_id),
do: :"streams_writer_worker_#{store}_#{stream_id}"
defp get_writer(store, stream_id) do
case get_cluster_writer(store, stream_id) do
nil ->
start_writer(store, stream_id)
writer_pid ->
writer_pid
end
end
defp get_cluster_writer(store, stream_id) do
case Swarm.registered()
|> Enum.filter(fn {name, _} ->
match?({:streams_writer_worker, ^store, ^stream_id}, name)
end)
|> Enum.map(fn {_, pid} -> pid end) do
[] ->
nil
writers ->
writers
|> 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_writer(store, stream_id) do
partition = partition_for(store, stream_id)
partition_name = ExESDB.StoreNaming.partition_name(ExESDB.StreamsWriters, store)
case DynamicSupervisor.start_child(
{:via, PartitionSupervisor, {partition_name, partition}},
{ExESDB.StreamsWriterWorker, {store, stream_id, partition}}
) do
{:ok, pid} -> pid
{:error, {:already_started, pid}} -> pid
{:error, reason} -> raise "failed to start streams writer: #{inspect(reason)}"
end
end
end