Packages
ex_esdb
0.0.8-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/event_stream_writer.ex
defmodule ExESDB.EventStreamWriter do
@moduledoc false
alias ExESDB.EventStoreInfo, as: ESInfo
# @retry_attempts 3
# @retry_delay 500
defp handle_transaction_result({:ok, {:commit, result}}), do: {:ok, result}
defp handle_transaction_result({:ok, {:abort, reason}}), do: {:error, reason}
defp handle_transaction_result({:error, reason}), do: {:error, reason}
def append_events_tx(store, stream_id, events) do
store
|> :khepri.transaction(fn ->
actual_version =
store
|> ESInfo.get_version!(stream_id)
store
|> append_events(stream_id, events, actual_version)
end)
|> handle_transaction_result()
end
def append_events(store, stream_id, events, current_version) do
events
|> Enum.reduce(
current_version,
fn event, version ->
new_version = version + 1
padded_version = ExESDB.VersionFormatter.pad_version(new_version, 6)
now =
DateTime.utc_now()
created = now
created_epoch =
now
|> DateTime.to_unix(:microsecond)
recorded_event =
event
|> to_event_record(
stream_id,
new_version,
created,
created_epoch
)
store
|> :khepri.put!([:streams, stream_id, padded_version], recorded_event)
new_version
end
)
end
defp to_event_record(
%ExESDB.NewEvent{} = new_event,
stream_id,
version,
created,
created_epoch
),
do: %ExESDB.EventRecord{
event_stream_id: stream_id,
event_number: version,
event_id: new_event.event_id,
event_type: new_event.event_type,
data_content_type: new_event.data_content_type,
metadata_content_type: new_event.metadata_content_type,
data: new_event.data,
metadata: new_event.metadata,
created: created,
created_epoch: created_epoch
}
end