Packages
ex_esdb
0.0.9-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/producer.ex
defmodule ExESDB.Repl.Producer do
@moduledoc false
use GenServer
require Logger
alias ExESDB.Repl.EventGenerator, as: ESGen
alias ExESDB.StreamsHelper, as: SHelper
alias ExESDB.StreamsWriter, as: StrWriter
alias ExESDB.Themes, as: Themes
defp append(store, stream, nbr_of_events) do
version =
store |> SHelper.get_version!(stream)
events = ESGen.generate_events(version, nbr_of_events)
{:ok, new_version} =
store
|> StrWriter.append_events(stream, version, events)
IO.puts("Appended #{nbr_of_events} events to #{stream}. Now at version #{new_version}.")
end
@impl true
def handle_info(:produce, state) do
stream = Keyword.get(state, :stream)
store = Keyword.get(state, :store)
batch_size = Keyword.get(state, :batch_size, 1)
store
|> append(stream, batch_size)
Process.send_after(self(), :produce, 2_000)
{:noreply, state}
end
############## PLUMBING ##############
@impl true
def init(args) do
Process.send_after(self(), :produce, 2_000)
{:ok, args}
end
def start_link(args) do
greenhouse = Keyword.get(args, :stream, "greenhouse0")
GenServer.start_link(
__MODULE__,
args,
name: Module.concat(__MODULE__, greenhouse)
)
end
def start_producer(args) do
greenhouse =
args
|> Keyword.get(:stream, "greenhouse0")
store =
args
|> Keyword.get(:store, :reg_gh)
case start_link(args) do
{:ok, pid} ->
IO.puts("#{Themes.producer(pid)} for [#{inspect(store)}:#{greenhouse}] is UP!")
pid
{:error, {:already_started, pid}} ->
IO.puts("#{Themes.producer(pid)} for [#{inspect(store)}:#{greenhouse}] is UP!")
pid
{:error, reason} ->
raise "Failed to start producer. Reason: #{inspect(reason)}"
end
end
end