defmodule Exkad.Node do use GenServer require Logger import Supervisor.Spec alias Exkad.Node.Buckets defmodule Descriptor do defstruct id: :empty, loc: :empty end alias Descriptor, as: ND @alpha 3 @k 10 @request_timeout 2000 def start_link(opts) do GenServer.start_link(__MODULE__, opts) end def init(opts) do state = opts |> init_desc |> init_tx |> start_supervised {:ok, state} end defp init_tx(%{tx: _} = state), do: state defp init_tx(_) do raise ArgumentError, message: "You need to provide a tx module" end defp bind_to(:self), do: self defp bind_to(address), do: address defp init_desc(state) do id = state[:id] bind = bind_to(state[:bind]) state = Dict.drop(state, [:id, :bind]) Enum.into(state, %{desc: %ND{id: id, loc: bind}}) end defp start_supervised(%{desc: node} = state) do children = [ worker(Exkad.Node.Buckets, [%{node: node, parent: self}]), worker(Exkad.Store, [%{parent: self}]), ] {:ok, sup} = Supervisor.start_link(children, strategy: :one_for_one) Process.link(sup) state end defp add_seed(%{seed: nil} = state), do: state defp add_seed(%{seed: seed} = state) do Buckets.add(state.buckets, seed) state end ## # There's no seed, so we're the first in the network defp add_seed(state), do: state defp populate_kbuckets(state = %{desc: %ND{id: this_id}}) do # Ask the seed for our id, which will give back a bunch of stuff to populate {_, state} = find_k(this_id, state) state end ## # There's no seed, first in the network.. defp populate_kbuckets(state), do: state defp unpack(<>) do [a, b, c, d, e, f, g, h] end def id_to_bin(loc) do loc |> :erlang.term_to_binary |> :erlang.binary_to_list |> Enum.map(fn num -> unpack(<>) end) |> List.flatten end def k_closest_nodes(to_id, %{buckets: buckets}) do Buckets.k_closest(buckets, to_id) end # return true if `closer_nodes` contains any nodes # that are closer to `to_find_id` than `nodes` defp are_closer?(closer_nodes, nodes, to_find_id) do Enum.any?(closer_nodes, fn %ND{id: cnode} -> closer_dist = Buckets.distance_from(cnode, to_find_id) Enum.any?(nodes, fn %ND{id: node} -> closer_dist < Buckets.distance_from(node, to_find_id) end) end) end # # Generic fan out forward to nodes # runs the to_do/1 function on each node # and waits for all nodes to return a result, # giving back a list of results defp forward(nodes, to_do) do nodes |> Enum.map(fn(node) -> {node, Task.async(fn -> to_do.(node) end)} end) |> Enum.map(fn {_, task} -> Task.await(task, @request_timeout) end) end # Helper for find_k/2 defp find_k(to_find_id, %{desc: desc, tx: tx} = state, nodes, seen) do closer = nodes |> forward(fn node -> Logger.debug("Finding #{to_find_id}") {:ok, node} = tx.find(node, to_find_id, desc) node end) |> List.flatten |> Enum.uniq |> Enum.reject(fn node -> node in seen or node == desc end) |> Buckets.k_closest_nodes_to(to_find_id) Buckets.add_all(state.buckets, closer) if are_closer?(closer, nodes, to_find_id) do find_k(to_find_id, state, closer, seen ++ nodes) else nodes = Buckets.k_closest_nodes_to(seen ++ nodes ++ closer, to_find_id) {nodes, state} end end # # Recursively find the k nodes that are closest to `to_find_id` # defp find_k(to_find_id, state) do find_k(to_find_id, state, k_closest_nodes(to_find_id, state), []) end defp extract_value([]), do: {:error, :not_found} defp extract_value([{:value, thing} | _]), do: {:value, thing} defp extract_value([_ | rest]), do: extract_value(rest) defp fan_out(nodes, fwd_fun, key, %{tx: tx}) do forward(nodes, fn node -> fwd_fun.(tx, node, key) end) |> List.flatten |> Enum.uniq end # # Forward a find request to `nodes` for key # defp get_from_many(nodes, key, state) do fan_out(nodes, fn(tx, node, key) -> tx.get(node, key) end, key, state) end defp search_from_many(nodes, key, state) do fan_out(nodes, fn(tx, node, key) -> tx.search(node, key) end, key, state) end defp is_search_result?({:ok, results}) when is_list(results), do: True defp is_search_result?(_), do: False defp flatten_search(results) do Enum.map(results, fn {:ok, res} -> res end) |> List.flatten end defp rsearch(key, %{desc: desc} = state, nodes, seen, results) do {search_results, closer} = search_from_many(nodes, key, state) |> Enum.partition(fn item -> is_search_result?(item) end) closer = closer |> Enum.reject(fn node -> node in seen or node == desc end) |> Buckets.k_closest_nodes_to(key) results = results ++ flatten_search(search_results) if are_closer?(closer, nodes, key) do rsearch(key, state, closer, seen ++ nodes, results) else closer_results = search_from_many(closer, key, state) |> Enum.filter(fn res -> is_search_result?(res) end) |> flatten_search closer_results ++ results end end # # Recursive find of a key # defp rfind_val(key, %{desc: desc} = state, nodes, seen) do closer = get_from_many(nodes, key, state) case extract_value(closer) do {:error, :not_found} -> Buckets.add_all(state.buckets, closer) closer = closer |> Enum.reject(fn node -> node in seen or node == desc end) |> Buckets.k_closest_nodes_to(key) if are_closer?(closer, nodes, key) do rfind_val(key, state, closer, seen ++ nodes) else val = get_from_many(closer, key, state) |> extract_value {state, val} end value -> {state, value} end end defp rfind_val(key, state) do case value_or_closest(key, state) do {:value, value} -> {state, {:value, value}} k_closest -> rfind_val(key, state, k_closest, []) end end defp rsearch(key, state) do nodes = k_closest_nodes(key, state) rsearch(key, state, nodes, [], []) end defp value_or_closest(key, %{store: store} = state) do case Exkad.Store.get(store, key) do {:error, :not_found} -> k_closest_nodes(key, state) value -> {:value, value} end end # # Replicate a keyval pair to k nodes # defp put_to_many(key, value, %{tx: tx} = state) do {nodes, _state} = find_k(key, state) forward(nodes, fn(node) -> tx.put(node, value) end) end # # Replicate metadata to k nodes # defp put_meta_to_many(meta_key, meta_term, ptr, %{tx: tx} = state) do {nodes, _state} = find_k(meta_key, state) forward(nodes, fn(node) -> tx.put_meta(node, meta_term, ptr) end) end def handle_call({:ping, from_desc}, _from, state) do Buckets.add(state.buckets, from_desc) {:reply, :ok, state} end def handle_call({:make_ping, to_desc}, _from, %{tx: tx, desc: from} = state) do res = tx.ping(to_desc, from) {:reply, res, state} end def handle_call(:dump, _from, %{buckets: buckets} = state) do buckets = Buckets.dump(buckets) {:reply, Enum.into(%{buckets: buckets}, state), state} end ## # This call is cheating, it will not be implemented over http, just erlang messaging # to make tests easier def handle_call(:describe, _from, %{desc: desc} = state) do {:reply, {:ok, desc}, state} end def handle_call({:find, to_find_id, from_desc}, _from, state) do closest = k_closest_nodes(to_find_id, state) Buckets.add(state.buckets, from_desc) {:reply, {:ok, closest}, state} end def handle_call({:put, value}, _from, %{store: store} = state) do result = Exkad.Store.put(store, value) {:reply, result, state} end def handle_call({:put_meta, term, ptr}, _from, %{store: store} = state) do result = Exkad.Store.put_meta(store, term, ptr) {:reply, result, state} end def handle_call({:search, key}, _from, %{store: store} = state) do result = Exkad.Store.search(store, key) {:reply, result, state} end def handle_call({:rsearch, term}, _from, state) do key = Exkad.Store.hash(term) result = rsearch(key, state) {:reply, result, state} end def handle_call({:rput, value, meta}, _from, %{store: store} = state) do result = Exkad.Store.put(store, value) {:ok, key} = result put_to_many(key, value, state) Enum.map(meta, fn {_meta_name, meta_term} -> {:ok, meta_key} = Exkad.Store.put_meta(store, meta_term, key) put_meta_to_many(meta_key, meta_term, key, state) end) {:reply, result, state} end def handle_call({:get, key}, _from, state) do result = value_or_closest(key, state) {:reply, result, state} end def handle_call({:rget, key}, _from, state) do {state, value} = rfind_val(key, state) {:reply, value, state} end ### # defp is_ready?(state) do Enum.all?([:store, :buckets], fn svc -> Dict.has_key?(state, svc) end) end defp register_service(:buckets, buckets, state) do state |> Dict.put(:buckets, buckets) |> add_seed |> populate_kbuckets end defp register_service(:store, store, state) do Enum.into(%{store: store}, state) end def handle_cast({:register, service, service_pid}, state) do state = register_service(service, service_pid, state) if is_ready?(state) do send(state.owner, :ready) end {:noreply, state} end def register(node_pid, service, service_pid) do GenServer.cast(node_pid, {:register, service, service_pid}) end end