Current section

Files

Jump to
horde lib horde registry_impl.ex
Raw

lib/horde/registry_impl.ex

defmodule Horde.RegistryImpl do
@moduledoc false
use GenServer
require Logger
defmodule State do
@moduledoc false
defstruct name: nil,
nodes: MapSet.new(),
members: MapSet.new(),
processes_updated_counter: 0,
processes_updated_at: 0,
registry_ets_table: nil,
pids_ets_table: nil,
keys_ets_table: nil
end
@spec child_spec(options :: list()) :: Supervisor.child_spec()
def child_spec(options \\ []) do
%{
id: Keyword.get(options, :name, __MODULE__),
start: {__MODULE__, :start_link, [options]}
}
end
@spec start_link(options :: list()) :: GenServer.on_start()
def start_link(options \\ []) do
name = Keyword.get(options, :name)
if !is_atom(name) || is_nil(name) do
raise ArgumentError, "expected :name to be given and to be an atom, got: #{inspect(name)}"
end
GenServer.start_link(__MODULE__, options, name: name)
end
### GenServer callbacks
def init(opts) do
{:ok, opts} =
case Keyword.get(opts, :init_module) do
nil -> {:ok, opts}
module -> module.init(opts)
end
Process.flag(:trap_exit, true)
name = Keyword.get(opts, :name)
pids_name = :"pids_#{name}"
keys_name = :"keys_#{name}"
Logger.info("Starting #{inspect(__MODULE__)} with name #{inspect(name)}")
unless is_atom(name) do
raise ArgumentError, "expected :name to be given and to be an atom, got: #{inspect(name)}"
end
:ets.new(name, [:named_table, {:read_concurrency, true}])
:ets.new(pids_name, [:named_table, {:read_concurrency, true}])
:ets.new(keys_name, [:named_table, {:read_concurrency, true}])
state = %State{
name: name,
registry_ets_table: name,
pids_ets_table: pids_name,
keys_ets_table: keys_name
}
state =
case Keyword.get(opts, :members) do
nil ->
state
members ->
members = Enum.map(members, &fully_qualified_name/1)
Enum.each(members, fn member ->
DeltaCrdt.mutate_async(crdt_name(state.name), :add, [{:member, member}, 1])
end)
neighbours = members -- [fully_qualified_name(state.name)]
send(crdt_name(state.name), {:set_neighbours, crdt_names(neighbours)})
%{state | nodes: Enum.map(members, fn {_name, node} -> node end) |> MapSet.new()}
end
case Keyword.get(opts, :meta) do
nil ->
nil
meta ->
Enum.each(meta, fn {key, value} -> put_meta(state, key, value) end)
end
{:ok, state}
end
def handle_info({:crdt_update, diffs}, state) do
new_state = process_diffs(state, diffs)
{:noreply, new_state}
end
def handle_info({:EXIT, pid, _reason}, state) do
case :ets.take(state.pids_ets_table, pid) do
[{_pid, keys}] ->
Enum.each(keys, fn key ->
DeltaCrdt.mutate_async(crdt_name(state.name), :remove, [{:key, key}])
:ets.match_delete(state.keys_ets_table, {key, {pid, :_}})
end)
_ ->
nil
end
{:noreply, state}
end
defp process_diffs(state, [diff | diffs]) do
process_diff(state, diff)
|> process_diffs(diffs)
end
defp process_diffs(state, []), do: state
defp process_diff(state, {:add, {:member, member}, 1}) do
new_members = MapSet.put(state.members, member)
send(
crdt_name(state.name),
{:set_neighbours, crdt_names(MapSet.delete(new_members, fully_qualified_name(state.name)))}
)
new_nodes = Enum.map(new_members, fn {_name, node} -> node end) |> MapSet.new()
%{state | members: new_members, nodes: new_nodes}
end
defp process_diff(state, {:remove, {:member, member}}) do
new_members = MapSet.delete(state.members, member)
new_nodes = Enum.map(new_members, fn {_name, node} -> node end) |> MapSet.new()
%{state | members: new_members, nodes: new_nodes}
end
defp process_diff(state, {:add, {:key, key}, {pid, value}}) do
link_local_pid(pid)
case :ets.lookup(state.pids_ets_table, pid) do
[] -> :ets.insert(state.pids_ets_table, {pid, [key]})
[{_pid, matches}] -> :ets.insert(state.pids_ets_table, {pid, Enum.uniq([key | matches])})
end
:ets.insert(state.keys_ets_table, {key, {pid, value}})
state
end
defp process_diff(state, {:remove, {:key, key}}) do
case :ets.lookup(state.keys_ets_table, key) do
[] ->
nil
[{key, {pid, _val}}] ->
case :ets.lookup(state.pids_ets_table, pid) do
[] -> []
[{pid, keys}] -> :ets.insert(state.pids_ets_table, {pid, List.delete(keys, key)})
end
end
:ets.match_delete(state.keys_ets_table, {key, :_})
state
end
defp process_diff(state, {:add, {:registry, key}, value}) do
:ets.insert(state.registry_ets_table, {key, value})
state
end
defp link_local_pid(pid) when node(pid) == node() do
Process.link(pid)
end
defp link_local_pid(_pid), do: nil
def handle_call({:set_members, members}, _from, state) do
new_members = MapSet.new(member_names(members))
Enum.each(MapSet.difference(state.members, new_members), fn removed_member ->
DeltaCrdt.mutate_async(crdt_name(state.name), :remove, [{:member, removed_member}])
end)
Enum.each(MapSet.difference(new_members, state.members), fn added_member ->
DeltaCrdt.mutate_async(crdt_name(state.name), :add, [{:member, added_member}, 1])
end)
neighbours = MapSet.difference(new_members, MapSet.new([state.name]))
send(crdt_name(state.name), {:set_neighbours, crdt_names(neighbours)})
{:reply, :ok, %{state | members: new_members}}
end
def handle_call({:register, key, value, pid}, _from, state) do
Process.link(pid)
DeltaCrdt.mutate_async(crdt_name(state.name), :add, [{:key, key}, {pid, value}])
case :ets.lookup(state.pids_ets_table, pid) do
[] ->
:ets.insert(state.pids_ets_table, {pid, [key]})
[{_pid, keys}] ->
:ets.insert(state.pids_ets_table, {pid, [key | keys]})
end
:ets.insert(state.keys_ets_table, {key, {pid, value}})
{:reply, {:ok, self()}, state}
end
def handle_call({:update_value, key, pid, value}, _from, state) do
DeltaCrdt.mutate_async(crdt_name(state.name), :add, [{:key, key}, {pid, value}])
:ets.insert(state.keys_ets_table, {key, {pid, value}})
{:reply, :ok, state}
end
def handle_call({:unregister, key, pid}, _from, state) do
DeltaCrdt.mutate_async(crdt_name(state.name), :remove, [{:key, key}])
case :ets.lookup(state.pids_ets_table, pid) do
[] -> []
[{pid, keys}] -> :ets.insert(state.pids_ets_table, {pid, List.delete(keys, key)})
end
:ets.match_delete(state.keys_ets_table, {key, {pid, :_}})
{:reply, :ok, state}
end
def handle_call({:put_meta, key, value}, _from, state) do
put_meta(state, key, value)
{:reply, :ok, state}
end
def handle_call(:members, _from, state) do
{:reply, {:ok, MapSet.to_list(state.members)}, state}
end
defp member_names(names) do
Enum.map(names, fn
{name, node} -> {name, node}
name when is_atom(name) -> {name, node()}
end)
end
defp crdt_names(names) do
Enum.map(names, fn {name, node} -> {crdt_name(name), node} end)
end
defp crdt_name(name), do: :"#{name}.Crdt"
defp fully_qualified_name({name, node}) when is_atom(name) and is_atom(node), do: {name, node}
defp fully_qualified_name(name) when is_atom(name), do: {name, node()}
defp put_meta(state, key, value) do
DeltaCrdt.mutate_async(crdt_name(state.name), :add, [{:registry, key}, value])
:ets.insert(state.registry_ets_table, {key, value})
end
end