Packages
x3m_system
0.7.2
0.9.1
0.9.0
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
retired
0.7.20
0.7.19
0.7.18
0.7.17
0.7.16
0.7.15
0.7.14
0.7.13
0.7.12
0.7.11
0.7.10
0.7.9
0.7.8
retired
0.7.7
0.7.6
retired
0.7.5
0.7.4
retired
0.7.3
retired
0.7.2
0.7.1
0.7.0
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
retired
0.5.6
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.9
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.1.1
0.1.0
Building blocks for distributed and/or CQRS/ES systems
Current section
Files
Jump to
Current section
Files
lib/aggregate_pid_facade.ex
defmodule X3m.System.AggregatePidFacade do
use GenServer
require Logger
alias X3m.System.AggregateSup
alias X3m.System.AggregateRegistry, as: Registry
def name(aggregate_mod),
do: Module.concat(aggregate_mod, PidFacade)
def get_aggregate_mod,
do: X3m.System.GenAggregate
## Client API
def start_link(aggregate_mod),
do: GenServer.start_link(__MODULE__, aggregate_mod, name: name(aggregate_mod))
def get_pid(server, key, when_not_registered),
do: GenServer.call(server, {:get_pid, key, when_not_registered})
def spawn_new(server, key, opts \\ []),
do: GenServer.call(server, {:spawn_new, key, opts})
def exit_process(server, key, reason),
do: GenServer.cast(server, {:exit_process, key, reason})
@impl GenServer
def init(aggregate_mod) do
state = %{
registry: Registry.name(aggregate_mod),
aggregate_sup: AggregateSup.name(aggregate_mod),
aggregate_mod: aggregate_mod
}
{:ok, state}
end
@impl GenServer
def handle_call({:spawn_new, key, opts}, _from, state) do
response = _spawn_new(key, state, opts)
{:reply, response, state}
end
def handle_call({:get_pid, key, when_not_registered}, _from, state),
do: {:reply, _get_pid(key, state, when_not_registered, []), state}
@impl GenServer
def handle_cast({:exit_process, key, reason}, state) do
case Registry.get(state.registry, key) do
{:ok, pid} ->
Logger.debug("Killing process #{inspect(pid)}: #{inspect(reason)}")
AggregateSup.terminate_child(state.aggregate_sup, pid)
_ ->
:ok
end
{:noreply, state}
end
defp _get_pid(key, state, when_not_registered, opts) do
case Registry.get(state.registry, key) do
:error ->
when_not_registered.(state.aggregate_mod, key, fn ->
_spawn_new(key, state, opts)
end)
{:ok, pid} ->
{:ok, pid}
end
end
defp _spawn_new(key, state, opts) do
{:ok, pid} = AggregateSup.start_child(state.aggregate_sup, opts)
X3m.System.Instrumenter.execute(:new_aggr_spawned, %{}, %{id: key})
case Registry.register(state.registry, key, pid) do
:ok -> {:ok, pid}
other -> other
end
end
end