diff --git a/README.md b/README.md index bdf738d..b877b48 100644 --- a/README.md +++ b/README.md @@ -29,10 +29,14 @@ In order to connect multiple nodes you should also set up [libcluster](https://g ## Configuration -Configure Squabble when you start the worker in your supervision tree. This should go _after_ `libcluster` if you're using that. All nodes should be connected before starting Squabble. +Configure Squabble when you start the worker in your supervision tree. This should go _after_ `:pg` and `libcluster` if you're using that. All nodes should be connected before starting Squabble. ```elixir children = [ + %{ + id: :pg, + start: {:pg, :start_link, []} + }, {Squabble, [subscriptions: [MyApp.Leader], size: 1]} ] ``` @@ -52,5 +56,13 @@ defmodule MyApp.Leader do @impl true def node_down() do end + + @impl true + def startup() do + end + + @impl true + def not_leader() do + end end ``` diff --git a/lib/squabble.ex b/lib/squabble.ex index 3db9d6f..16d1866 100644 --- a/lib/squabble.ex +++ b/lib/squabble.ex @@ -90,6 +90,10 @@ defmodule Squabble do size = Keyword.get(opts, :size, 1) subscriptions = Keyword.get(opts, :subscriptions, []) + Enum.each(subscriptions, fn module -> + Code.ensure_loaded(module) + end) + state = %State{ state: "candidate", size: size, diff --git a/lib/squabble/leader.ex b/lib/squabble/leader.ex index 394d366..e316249 100644 --- a/lib/squabble/leader.ex +++ b/lib/squabble/leader.ex @@ -14,4 +14,14 @@ defmodule Squabble.Leader do A node went down, callback from the squabble leader """ @callback node_down() :: :ok + + @doc """ + Callback if we're not the leader + """ + @callback not_leader(election_term()) :: :ok + + @doc""" + Callback on startup + """ + @callback startup(election_term()) :: :ok end diff --git a/lib/squabble/pg.ex b/lib/squabble/pg.ex index d03177b..b5a8df8 100644 --- a/lib/squabble/pg.ex +++ b/lib/squabble/pg.ex @@ -10,8 +10,7 @@ defmodule Squabble.PG do """ @spec join() :: :ok def join() do - :ok = :pg2.create(@key) - :ok = :pg2.join(@key, self()) + :ok = :pg.join(@key, self()) end @doc """ @@ -35,7 +34,7 @@ defmodule Squabble.PG do """ @spec members() :: [pid()] def members() do - :pg2.get_members(:squabble) + :pg.get_members(:squabble) end @doc """ diff --git a/lib/squabble/server.ex b/lib/squabble/server.ex index b658657..bad4cb2 100644 --- a/lib/squabble/server.ex +++ b/lib/squabble/server.ex @@ -56,6 +56,12 @@ defmodule Squabble.Server do Squabble.leader_check(pid) end) + Enum.each(winner_subscriptions(state), fn module -> + if function_exported?(module, :startup, 1) do + module.startup(state.term) + end + end) + {:ok, state} end @@ -82,9 +88,10 @@ defmodule Squabble.Server do "Starting an election for term #{term}, announcing candidacy" end, type: :squabble) + case check_term_newer(state, term) do {:ok, :newer} -> - if state.size == 1 do + if length(PG.members()) == 1 do voted_leader(state, 1) else PG.broadcast(fn pid -> @@ -105,6 +112,12 @@ defmodule Squabble.Server do "Someone already won this round, not starting" end, type: :squabble) + Enum.each(winner_subscriptions(state), fn module -> + if function_exported?(module, :not_leader, 1) do + module.not_leader(state.term) + end + end) + {:ok, state} {:error, :older} -> @@ -112,6 +125,14 @@ defmodule Squabble.Server do "This term has already completed, not starting" end, type: :squabble) + if state.state != "leader" do + Enum.each(winner_subscriptions(state), fn module -> + if function_exported?(module, :not_leader, 1) do + module.not_leader(state.term) + end + end) + end + {:ok, state} end end @@ -187,6 +208,14 @@ defmodule Squabble.Server do :ets.insert(@key, {:is_leader?, false}) + if leader_node != node() do + Enum.each(winner_subscriptions(state), fn module -> + if function_exported?(module, :not_leader, 1) do + module.not_leader(term) + end + end) + end + state = state |> Map.put(:term, term) @@ -305,7 +334,11 @@ defmodule Squabble.Server do """ @spec check_majority_votes(State.t()) :: {:ok, :majority} | {:error, :not_enough} def check_majority_votes(state) do - case length(state.votes) >= state.size / 2 do + + current_size = length(PG.members()) + current_size = if current_size < 1, do: 1, else: current_size + + case length(state.votes) >= current_size / 2 do true -> {:ok, :majority} @@ -380,4 +413,13 @@ defmodule Squabble.Server do {:ok, state} end end + + @impl true + def not_leader(_term) do + :ok + end + @impl true + def startup(_term) do + :ok + end end