Packages
ex_esdb
0.0.11-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 to interact with the ExESDB.system,
running a store called "reg_gh" (Regulate Greenhouse)
"""
alias ExESDB.Repl.EventGenerator, as: ESGen
alias ExESDB.Repl.Observer, as: Observer
alias ExESDB.Repl.Producer, as: Producer
alias ExESDB.Repl.Subscriber, as: Subscriber
alias ExESDB.GatewayAPI, as: API
require Logger
@store :reg_gh
@greenhouse1 "greenhouse1"
@greenhouse2 "greenhouse2"
@greenhouse3 "greenhouse3"
@greenhouse4 "greenhouse4"
@greenhouse5 "greenhouse5"
@greenhouses [
@greenhouse1,
@greenhouse2,
@greenhouse3,
@greenhouse4,
@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_opts, do: ExESDB.Options.app_env()
def get_streams, do: API.get_streams(@store)
def get_subscriptions, do: API.get_subscriptions(@store)
@doc """
Append events to a stream.
"""
@spec append(
stream :: atom(),
nbr_of_events :: integer()
) :: {:ok, list(), integer()} | {:error, term()}
def append(stream, nbr_of_events) do
{:ok, version} = API.get_version(@store, stream)
events = ESGen.generate_events(version, nbr_of_events)
case @store
|> API.append_events(stream, events) do
{:ok, new_version} ->
{:ok, result} =
@store
|> API.get_events(stream, 1, new_version, :forward)
{:ok, result, result |> Enum.count()}
{:error, reason} ->
{:error, reason}
end
end
def start_observer_for_all_streams do
Observer.start(store: @store, type: :by_stream, selector: "$all")
end
def start_greenhouse1_subscriber do
Subscriber.start(store: @store, type: :by_stream, selector: "$greenhouse1")
end
def start_producers do
ESGen.streams()
|> Enum.each(fn stream -> Producer.start(stream) end)
end
end