Packages
exq
0.16.0
0.24.0
0.23.0
0.22.0
0.21.0
0.20.0
0.19.0
0.18.0
0.17.0
0.16.2
0.16.1
0.16.0
0.15.0
0.14.0
0.13.5
0.13.4
0.13.3
0.13.2
0.13.1
0.13.0
0.12.2
0.12.1
0.12.0
0.11.0
0.10.1
0.10.0
0.9.1
0.9.0
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.3
0.7.2
0.7.1
0.7.0
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.0
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.0
0.2.3
0.2.2
0.2.1
0.2.0
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
0.0.2
Exq is a job processing library compatible with Resque / Sidekiq for the Elixir language.
Current section
Files
Jump to
Current section
Files
lib/exq/heartbeat/monitor.ex
defmodule Exq.Heartbeat.Monitor do
use GenServer
require Logger
alias Exq.Redis.Heartbeat
alias Exq.Support.Config
defmodule State do
defstruct [
:namespace,
:redis,
:interval,
:queues,
:node_id,
:missed_heartbeats_allowed,
:stats
]
end
def start_link(options) do
GenServer.start_link(
__MODULE__,
%State{
namespace: Keyword.fetch!(options, :namespace),
redis: Keyword.fetch!(options, :redis),
interval: Keyword.fetch!(options, :heartbeat_interval),
node_id: Keyword.get(options, :node_id, Config.node_identifier().node_id()),
queues: Keyword.fetch!(options, :queues),
stats: Keyword.get(options, :stats),
missed_heartbeats_allowed: Keyword.fetch!(options, :missed_heartbeats_allowed)
},
[]
)
end
def init(state) do
:ok = schedule_verify(state.interval)
{:ok, state}
end
def handle_info(:verify, state) do
case Heartbeat.dead_nodes(
state.redis,
state.namespace,
state.interval,
state.missed_heartbeats_allowed
) do
{:ok, node_ids} ->
Enum.each(node_ids, fn {node_id, score} ->
:ok = re_enqueue_backup(state, node_id, score)
end)
_error ->
:ok
end
:ok = schedule_verify(state.interval)
{:noreply, state}
end
def handle_info(msg, state) do
Logger.error("Received unexpected info message in #{__MODULE__} #{inspect(msg)}")
{:noreply, state}
end
defp schedule_verify(interval) do
_reference = Process.send_after(self(), :verify, interval)
:ok
end
defp re_enqueue_backup(state, node_id, score) do
Logger.info(
"#{node_id} missed the last #{state.missed_heartbeats_allowed} heartbeats. Re-enqueing jobs from backup and cleaning up stats."
)
Enum.uniq(Exq.Redis.JobQueue.list_queues(state.redis, state.namespace) ++ state.queues)
|> Enum.each(fn queue ->
Heartbeat.re_enqueue_backup(state.redis, state.namespace, node_id, queue, score)
end)
if state.stats do
:ok = Exq.Stats.Server.cleanup_host_stats(state.stats, state.namespace, node_id)
end
_ = Heartbeat.unregister(state.redis, state.namespace, node_id)
:ok
end
end