Packages
ex_esdb
0.0.6-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_reader.ex
defmodule ExESDB.EventStreamReader do
@moduledoc false
import ExESDB.Khepri.Conditions
def get_current_version!(store, stream_id) do
case store
|> get_current_version(stream_id) do
{:ok, count} -> count
_ -> 0
end
end
def get_current_version(store, stream_id),
do:
store
|> :khepri.count([
:streams,
stream_id,
if_node_exists(exists: true)
])
def read_events(store, stream_id, start_version, count) do
start_version..(start_version + count - 1)
|> Enum.map(fn version ->
padded_version = ExESDB.VersionFormatter.pad_version(version, 6)
store
|> :khepri.get!([:streams, stream_id, padded_version])
end)
|> Enum.reject(&is_nil/1)
end
def get_streams(store) do
store
|> :khepri.get!([:streams])
|> Enum.reduce([], fn {stream_id, _stream}, acc -> stream_id ++ acc end)
end
end