Packages
extreme
0.5.1
1.1.4
1.1.3
1.1.2
1.1.1
1.1.1-rc01
1.1.0
1.1.0-rc9
1.1.0-rc8
1.1.0-rc7
1.1.0-rc6
1.1.0-rc5
1.1.0-rc4
1.1.0-rc3
1.1.0-rc2
1.1.0-rc1
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.13.4
0.13.3
0.13.2
0.13.1
0.13.0
0.12.1
0.12.0
0.11.0
0.10.4
0.10.3
0.10.2
0.10.1
0.10.0
0.9.2
0.9.1
0.9.0
0.8.1
0.8.0
0.7.1
0.7.0
0.6.2
0.6.1
0.6.0
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.3
0.4.2
0.4.1
Elixir TCP client for EventStore.
Current section
Files
Jump to
Current section
Files
lib/cluster_connection.ex
defmodule Extreme.ClusterConnection do
require Logger
def get_node(connection_settings) do
nodes = Keyword.fetch! connection_settings, :nodes
gossip_timeout = Keyword.get connection_settings, :gossip_timeout, 1_000
gossip_with nodes, gossip_timeout
end
def get_node(:cluster_dns, connection_settings) do
hostname = Keyword.fetch! connection_settings, :host
hostname = if is_binary(hostname), do: to_char_list(hostname)
{:ok, ips}= :inet.getaddrs(hostname, :inet, 1_000)
gossip_timeout = Keyword.get connection_settings, :gossip_timeout, 1_000
gossip_port= Keyword.get connection_settings, :port, 2113
nodes= Enum.map ips, fn ip -> %{host: to_string(:inet.ntoa ip), port: gossip_port} end
gossip_with nodes, gossip_timeout
end
defp gossip_with([], _), do: {:error, :no_more_gossip_seeds}
defp gossip_with([node|rest_nodes], gossip_timeout) do
url = "http://#{node.host}:#{node.port}/gossip?format=json"
Logger.info "Gossip with #{url}"
case HTTPoison.get url, [], timeout: gossip_timeout do
{:ok, %HTTPoison.Response{status_code: 200, body: body}} ->
Poison.decode!(body)
|> choose_node
error ->
Logger.error "Error getting gossip: #{inspect error}"
gossip_with rest_nodes, gossip_timeout
end
end
defp choose_node(%{"members" => members}) do
best_candidate = members
|> get_alive
|> inject_state_rank
|> remove_0_ranks
|> sort_by_rank
|> List.first
Logger.info "We've chosen node: #{inspect best_candidate}"
{:ok, best_candidate["externalTcpIp"], best_candidate["externalTcpPort"]}
end
defp get_alive(members), do: Enum.filter(members, fn(m) -> m["isAlive"] end)
defp inject_state_rank(members), do: Enum.map(members, fn(m) -> Dict.merge m, %{state_rank: rank_state(m["state"])} end)
defp remove_0_ranks(members), do: Enum.reject(members, &(&1.state_rank == 0))
defp sort_by_rank(candidates), do: Enum.sort(candidates, &(&1.state_rank < &2.state_rank))
defp rank_state("Master"), do: 1
defp rank_state("PreMaster"), do: 2
defp rank_state("Slave"), do: 3
defp rank_state("Clone"), do: 4
defp rank_state("CatchingUp"), do: 5
defp rank_state("PreReplica"), do: 6
defp rank_state("Unknown"), do: 7
defp rank_state("Initializing"), do: 8
defp rank_state("Manager"), do: 0
defp rank_state("ShuttingDown"), do: 0
defp rank_state("Shutdown"), do: 0
defp rank_state(state), do: Logger.warn("Unrecognized node state: #{state}"); 0
end