defmodule Cluster.EcsStrategy do @moduledoc """ Clustering!! """ use GenServer use Cluster.Strategy alias Cluster.Strategy.State @default_polling_interval 5_000 @impl true def start_link(args), do: GenServer.start_link(__MODULE__, args) @impl true def init([%State{meta: nil} = state]) do {:ok, _pid} = Cluster.EcsClusterInfo.start_link(state.config) init([%State{state | :meta => MapSet.new()}]) end def init([%State{} = state]) do {:ok, load(state), 0} end @impl true def handle_info(:load, state) do {:noreply, load(state)} end def handle_info(_, state) do {:noreply, state} end defp load(%State{topology: topology, meta: meta} = state) do new_nodelist = MapSet.new(get_nodes(state)) removed = MapSet.difference(meta, new_nodelist) new_nodelist = case Cluster.Strategy.disconnect_nodes( topology, state.disconnect, state.list_nodes, MapSet.to_list(removed) ) do :ok -> new_nodelist {:error, bad_nodes} -> # Add back the nodes which should have been removed, but which couldn't be for some reason Enum.reduce(bad_nodes, new_nodelist, fn {n, _}, acc -> MapSet.put(acc, n) end) end new_nodelist = case Cluster.Strategy.connect_nodes( topology, state.connect, state.list_nodes, MapSet.to_list(new_nodelist) ) do :ok -> new_nodelist {:error, bad_nodes} -> # Remove the nodes which should have been added, but couldn't be for some reason Enum.reduce(bad_nodes, new_nodelist, fn {n, _}, acc -> MapSet.delete(acc, n) end) end Process.send_after( self(), :load, polling_interval(state) ) %State{state | meta: new_nodelist} end @spec get_nodes(State.t()) :: [atom()] defp get_nodes(%State{topology: _topology, config: _config}) do Cluster.EcsClusterInfo.get_nodes() |> Map.keys() end defp polling_interval(%State{config: config}) do Keyword.get(config, :polling_interval, @default_polling_interval) end end