Packages
ex_esdb
0.7.0
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/emitter_pool.ex
defmodule ExESDB.EmitterPool do
@moduledoc false
use Supervisor
require Logger
alias ExESDB.Themes, as: Themes
def name(store, sub_topic),
do: :"#{store}:#{sub_topic}_emitter_pool"
def start_link({store, sub_topic, subscriber, pool_size, filter}) do
Supervisor.start_link(
__MODULE__,
{store, sub_topic, subscriber, pool_size, filter},
name: name(store, sub_topic)
)
end
@impl Supervisor
def init({store, sub_topic, subscriber, pool_size, filter}) do
emitter_names =
store
|> :emitter_group.setup_emitter_mechanism(sub_topic, filter, pool_size)
children =
for emitter <- emitter_names do
Supervisor.child_spec(
{ExESDB.EmitterWorker, {store, sub_topic, subscriber, emitter}},
id: emitter
)
end
# Enhanced prominent multi-line pool startup message
pool_name = name(store, sub_topic)
emitter_count = length(emitter_names)
IO.puts("")
IO.puts("")
IO.puts("┏━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┓")
IO.puts("#{Themes.emitter_pool_success_msg(self(), " 🚀 EMITTER POOL STARTUP 🚀 ")}")
IO.puts("┣━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┫")
IO.puts("#{Themes.emitter_pool_success_msg(self(), " Pool Name: #{pool_name}")}")
IO.puts("#{Themes.emitter_pool_success_msg(self(), " Store ID: #{store}")}")
IO.puts("#{Themes.emitter_pool_success_msg(self(), " Topic: #{sub_topic}")}")
IO.puts("#{Themes.emitter_pool_success_msg(self(), " Workers: #{emitter_count}")}")
IO.puts("#{Themes.emitter_pool_success_msg(self(), " PID: #{inspect(self())}")}")
IO.puts("┗━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┛")
IO.puts("")
IO.puts("")
Supervisor.init(children, strategy: :one_for_one)
end
def stop(store, sub_topic) do
pool_name = name(store, sub_topic)
# Enhanced prominent multi-line pool shutdown message
IO.puts("")
IO.puts("")
IO.puts("┏━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┓")
IO.puts("#{Themes.emitter_pool_failure_msg(self(), " 🚨 EMITTER POOL SHUTDOWN 🚨 ")}")
IO.puts("┣━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┫")
IO.puts("#{Themes.emitter_pool_failure_msg(self(), " Pool Name: #{pool_name}")}")
IO.puts("#{Themes.emitter_pool_failure_msg(self(), " Store ID: #{store}")}")
IO.puts("#{Themes.emitter_pool_failure_msg(self(), " Topic: #{sub_topic}")}")
IO.puts("#{Themes.emitter_pool_failure_msg(self(), " Reason: Manual Stop")}")
IO.puts("┗━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━┛")
IO.puts("")
IO.puts("")
Supervisor.stop(pool_name)
end
end