Packages
exq
0.19.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
test/exq/heartbeat/monitor_test.exs
defmodule Exq.Heartbeat.MonitorTest do
use ExUnit.Case
import ExqTestUtil
alias Exq.Support.Config
alias Exq.Redis.Heartbeat
@opts [
redis: :testredis,
heartbeat_enable: true,
heartbeat_interval: 200,
missed_heartbeats_allowed: 3,
queues: ["default"],
namespace: "test",
name: ExqHeartbeat,
stats: ExqHeartbeat.Stats
]
setup do
TestRedis.setup()
on_exit(fn ->
wait()
TestRedis.teardown()
end)
end
test "re-enqueues orphaned jobs from dead node's backup queue" do
{:ok, _} = Exq.Stats.Server.start_link(@opts)
redis = :testredis
servers =
for i <- 1..5 do
{:ok, heartbeat} =
Exq.Heartbeat.Server.start_link(Keyword.put(@opts, :node_id, to_string(i)))
{:ok, monitor} =
Exq.Heartbeat.Monitor.start_link(Keyword.put(@opts, :node_id, to_string(i)))
%{heartbeat: heartbeat, monitor: monitor}
end
assert {:ok, 1} = working(redis, "3")
Process.sleep(1000)
assert alive_nodes(redis) == ["1", "2", "3", "4", "5"]
assert queue_length(redis, "3") == {:ok, 1}
server = Enum.at(servers, 2)
:ok = GenServer.stop(server.heartbeat)
Process.sleep(2000)
assert alive_nodes(redis) == ["1", "2", "4", "5"]
assert queue_length(redis, "3") == {:ok, 0}
end
test "re-enqueues more than 10 orphaned jobs from dead node's backup queue" do
{:ok, _} = Exq.Stats.Server.start_link(@opts)
redis = :testredis
servers =
for i <- 1..5 do
{:ok, heartbeat} =
Exq.Heartbeat.Server.start_link(Keyword.put(@opts, :node_id, to_string(i)))
{:ok, monitor} =
Exq.Heartbeat.Monitor.start_link(Keyword.put(@opts, :node_id, to_string(i)))
%{heartbeat: heartbeat, monitor: monitor}
end
for i <- 1..15 do
assert {:ok, ^i} = working(redis, "3")
end
Process.sleep(1000)
assert alive_nodes(redis) == ["1", "2", "3", "4", "5"]
assert queue_length(redis, "3") == {:ok, 15}
server = Enum.at(servers, 2)
:ok = GenServer.stop(server.heartbeat)
Process.sleep(2000)
assert alive_nodes(redis) == ["1", "2", "4", "5"]
assert queue_length(redis, "3") == {:ok, 0}
end
test "can handle connection failure" do
with_application_env(:exq, :redis_timeout, 500, fn ->
{:ok, _} = Exq.Stats.Server.start_link(@opts)
redis = :testredis
assert alive_nodes(redis) == []
{:ok, _} = Exq.Heartbeat.Server.start_link(Keyword.put(@opts, :node_id, "1"))
{:ok, _} = Exq.Heartbeat.Monitor.start_link(Keyword.put(@opts, :node_id, "1"))
spawn(fn ->
Redix.command(:testredis, ["DEBUG", "SLEEP", "2"])
end)
Process.sleep(2000)
Redix.command(:testredis, ["FLUSHALL"])
Process.sleep(1000)
assert alive_nodes(redis) == ["1"]
end)
end
test "shouldn't dequeue from live node" do
redis = :testredis
namespace = Config.get(:namespace)
interval = 100
missed_heartbeats_allowed = 3
Heartbeat.register(redis, namespace, "1")
assert {:ok, 1} = working(redis, "1")
Process.sleep(1000)
{:ok, %{"1" => score}} =
Heartbeat.dead_nodes(
redis,
namespace,
interval,
missed_heartbeats_allowed
)
assert queue_length(redis, "1") == {:ok, 1}
Heartbeat.re_enqueue_backup(redis, namespace, "1", "default", score)
assert queue_length(redis, "1") == {:ok, 0}
Heartbeat.register(redis, namespace, "1")
assert {:ok, 1} = working(redis, "1")
Process.sleep(1000)
{:ok, %{"1" => score}} =
Heartbeat.dead_nodes(
redis,
namespace,
interval,
missed_heartbeats_allowed
)
# The node came back after we got the dead node list, but before we could re-enqueue
Heartbeat.register(redis, namespace, "1")
Heartbeat.re_enqueue_backup(redis, namespace, "1", "default", score)
assert queue_length(redis, "1") == {:ok, 1}
Process.sleep(1000)
{:ok, %{"1" => score}} =
Heartbeat.dead_nodes(
redis,
namespace,
interval,
missed_heartbeats_allowed
)
# The node got removed by another heartbeat monitor
Heartbeat.unregister(redis, namespace, "1")
Heartbeat.re_enqueue_backup(redis, namespace, "1", "default", score)
assert queue_length(redis, "1") == {:ok, 1}
end
defp alive_nodes(redis) do
{:ok, nodes} =
Redix.command(redis, ["ZRANGEBYSCORE", "#{Config.get(:namespace)}:heartbeats", "0", "+inf"])
Enum.sort(nodes)
end
defp working(redis, node_id) do
Redix.command(redis, [
"LPUSH",
Exq.Redis.JobQueue.backup_queue_key(Config.get(:namespace), node_id, "default"),
"{}"
])
end
defp queue_length(redis, node_id) do
Redix.command(redis, [
"LLEN",
Exq.Redis.JobQueue.backup_queue_key(Config.get(:namespace), node_id, "default")
])
end
end