Packages
x3m_system
0.8.3
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_registry.ex
defmodule X3m.System.AggregateRegistry do
@moduledoc """
Keeps track of registered aggregate pids.
"""
require Logger
use GenServer
def name(aggregate_mod),
do: Module.concat(aggregate_mod, PidRegistry)
## Client API
@doc """
Starts process registry
"""
def start_link(name),
do: GenServer.start_link(__MODULE__, name, name: name)
@doc """
Returns `true` if `key` is already registered in `server`, `false` otherwise.
"""
def has_key?(server, key), do: GenServer.call(server, {:has_key?, key})
@doc """
Looks up the process pid for `id` stored in `server`.
Returns `{:ok, pid}` if the one exists, `:error` otherwise.
"""
def get(server, key),
do: server |> pid_table |> _get(key)
@doc """
Ensures there is a `pid` associated with `key` in `server`.
Returns :ok once when process is successfully registered
"""
def register(server, key, pid),
do: GenServer.call(server, {:register, key, pid})
## Server Callbacks
def init(name) do
name
|> pid_table
|> :ets.new([:named_table, read_concurrency: true])
name
|> ref_table
|> :ets.new([:named_table, read_concurrency: true])
{:ok, %{pid_table: pid_table(name), ref_table: ref_table(name)}}
end
## Private functions
defp _get(table, key) do
case :ets.lookup(table, key) do
[{^key, val}] -> {:ok, val}
[] -> :error
end
end
defp pid_table(name), do: Module.concat(name, Pids)
defp ref_table(name), do: Module.concat(name, Refs)
def handle_call({:has_key?, key}, _from, state) do
{:reply, Map.has_key?(state.processes, key), state}
end
def handle_call({:get, key}, _from, state) do
{:reply, Map.fetch(state.processes, key), state}
end
def handle_call({:register, key, pid}, _from, state) do
result =
case _get(state.pid_table, key) do
{:ok, pid} ->
{:error, :key_already_registered, pid}
:error ->
ref = Process.monitor(pid)
true = :ets.insert(state.pid_table, {key, pid})
true = :ets.insert(state.ref_table, {ref, key})
:ok
end
{:reply, result, state}
end
def handle_info({:DOWN, ref, :process, pid, reason}, state) do
case _get(state.ref_table, ref) do
{:ok, key} ->
Logger.warning("Process #{inspect(pid)} went down: #{inspect(reason)}")
:ets.delete(state.ref_table, ref)
:ets.delete(state.pid_table, key)
end
{:noreply, state}
end
def handle_info(_msg, state), do: {:noreply, state}
end