Packages
extreme
1.1.0
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/extreme/cluster_connection.ex
defmodule Extreme.ClusterConnection do
@moduledoc """
Finds node in EventStore cluster that connection should be established with.
"""
require Logger
def gossip_with(nodes, gossip_timeout, mode)
def gossip_with([], _, _), do: {:error, :no_more_gossip_seeds}
def gossip_with([node | rest_nodes], gossip_timeout, mode) do
url = ~c"http://#{node.host}:#{node.port}/gossip?format=json"
Logger.info("Gossip with #{url}")
case :httpc.request(:get, {url, []}, [timeout: gossip_timeout], []) do
{:ok, {{_version, 200, _status}, _headers, body}} ->
body
|> Jason.decode!()
|> _choose_node(mode)
error ->
Logger.error("Error getting gossip: #{inspect(error)}")
gossip_with(rest_nodes, gossip_timeout, mode)
end
end
defp _choose_node(%{"members" => members}, mode) do
best_candidate =
members
|> _get_alive
|> _inject_state_rank(mode)
|> _remove_0_ranks
|> _sort_by_rank
|> List.first()
Logger.info("We've chosen node: #{inspect(best_candidate)}")
{:ok, String.to_charlist(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, mode),
do: Enum.map(members, fn m -> Map.merge(m, %{state_rank: _rank_state(m["state"], mode)}) 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))
# Prefer Master when writing, but Slave for everything else (read)
defp _rank_state("Master", :write), do: 1
defp _rank_state("PreMaster", :write), do: 2
defp _rank_state("Slave", :write), do: 3
defp _rank_state("Slave", _), do: 1
defp _rank_state("Master", _), do: 2
defp _rank_state("PreMaster", _), 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.warning("Unrecognized node state: #{state}")
end