Packages
ex_esdb
0.0.5-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/repl.ex
defmodule ExESDB.Repl do
@moduledoc """
This module is used to start a REPL with the ExESDB.system running.
"""
alias ExESDB.Repl.EventGenerator, as: EventGenerator
alias ExESDB.Repl.EventStreamMonitor, as: EventStreamMonitor
alias ExESDB.EventStoreInfo, as: ESInfo
require Logger
@store :reg_gh
@greenhouse1 :greenhouse1
@greenhouse2 :greenhouse2
@greenhouse3 :greenhouse3
@greenhouse4 :greenhouse4
@greenhouse5 :greenhouse5
def store, do: @store
def stream1, do: @greenhouse1
def stream2, do: @greenhouse2
def stream3, do: @greenhouse3
def stream4, do: @greenhouse4
def stream5, do: @greenhouse5
def get_config, do: ExESDB.Options.app_env()
def get_streams, do: ESInfo.get_streams_raw(@store)
def start_monitor(store) do
case store
|> ExESDB.EventStore.get_state() do
{:ok, [config: config, store: _]} ->
IO.puts "Starting monitor for #{inspect(config, pretty: true)}"
{:ok, _pid} = EventStreamMonitor.start_link(config)
{:error, reason} -> raise "Failed to get state. Reason: #{inspect(reason)}"
end
end
def initialize(stream) do
initialized = EventGenerator.initialize()
{:ok, actual_version} =
@store
|> ExESDB.EventStore.stream_version(stream)
{:ok, new_version} =
@store
|> ExESDB.EventStore.append_to_stream(stream, actual_version, [initialized])
{:ok, result} =
@store
|> ExESDB.EventStore.read_stream_forward(stream, 1, new_version)
{:ok, result, result |> Enum.count()}
end
def append(stream, nbr_of_events) do
case stream |> all() do
nil ->
stream
|> initialize()
stream
|> append(nbr_of_events)
_ ->
events =
EventGenerator.generate_events(nbr_of_events)
{:ok, actual_version} =
@store
|> ExESDB.EventStore.stream_version(stream)
{:ok, new_version} =
@store
|> ExESDB.EventStore.append_to_stream(stream, actual_version, events)
{:ok, result} =
@store
|> ExESDB.EventStore.read_stream_forward(stream, 1, new_version)
{:ok, result, result |> Enum.count()}
end
end
def all(stream) do
case @store
|> ExESDB.EventStore.stream_version(stream) do
{:ok, 0} ->
nil
{:ok, version} ->
{:ok, events} =
@store
|> ExESDB.EventStore.read_stream_forward(stream, 1, version)
events
end
end
end