Packages
ex_esdb
0.0.16
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.Options, as: Options
alias ExESDB.Repl.EventGenerator, as: ESGen
alias ExESDB.StreamsHelper, as: SHelper
alias ExESDB.StreamsWriter, as: StrWriter
alias ExESDB.Themes, as: Themes
@impl true
def handle_info(:produce, state) do
stream_id = stream_id(state)
store = store(state)
batch_size = batch_size(state)
period = period(state)
store
|> append(stream_id, batch_size)
produce(period)
{:noreply, state}
end
############## PLUMBING ##############
@impl true
def init(args) do
period = period(args)
produce(period)
{:ok, args}
end
def start_link(args) do
stream_id = stream_id(args)
GenServer.start_link(
__MODULE__,
args,
name: Module.concat(__MODULE__, stream_id)
)
end
@doc """
Starts a producer process for testing purposes.
## Parameters
* `stream_id`: The id of the stream to append events to (string, default: `greenhouse0`).
* `batch_size`: The number of events to append in a single batch (integer, default: `1`).
* `period`: The time in milliseconds between appending events (integer, default: `2000`).
"""
@spec start(
stream_id :: String.t(),
batch_size :: integer(),
period :: integer()
) :: pid()
def start(stream_id \\ "greenhouse0", batch_size \\ 1, period \\ 2_000) do
store = Options.store_id()
args = [store: store, stream_id: stream_id, batch_size: batch_size, period: period]
case start_link(args) do
{:ok, pid} ->
IO.puts("#{Themes.producer(pid)} for [#{inspect(store)}:#{stream_id}] is UP!")
pid
{:error, {:already_started, pid}} ->
IO.puts("#{Themes.producer(pid)} for [#{inspect(store)}:#{stream_id}] is already UP!")
pid
{:error, reason} ->
raise "Failed to start producer. Reason: #{inspect(reason)}"
end
end
@spec stop(stream_id :: String.t()) :: :ok
def stop(stream_id),
do:
GenServer.stop(
Module.concat(__MODULE__, stream_id),
:normal
)
defp store(args), do: Keyword.get(args, :store, :reg_gh)
defp stream_id(args), do: Keyword.get(args, :stream_id, "greenhouse0")
defp batch_size(args), do: Keyword.get(args, :batch_size, 1)
defp period(args), do: Keyword.get(args, :period, 2_000)
defp produce(period), do: Process.send_after(self(), :produce, :rand.uniform(period * 5))
defp append(store, stream_id, nbr_of_events) do
version =
store |> SHelper.get_version!(stream_id)
events = ESGen.generate_events(version, nbr_of_events)
{:ok, new_version} =
store
|> StrWriter.append_events(stream_id, version, events)
msg = "#{nbr_of_events} event(s) to #{inspect(store)}:#{stream_id}.v(#{new_version})"
IO.puts(Themes.appended(msg))
end
end