Packages
exq
0.11.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/job_queue_test.exs
defmodule JobQueueTest do
use ExUnit.Case
alias Exq.Redis.JobQueue
alias Exq.Support.Job
alias Exq.Support.Time
import ExqTestUtil
@host 'host-name'
setup do
TestRedis.setup
on_exit fn ->
TestRedis.teardown
end
end
def assert_dequeue_job(queues, expected_result) do
jobs = JobQueue.dequeue(:testredis, "test", @host, queues)
result = jobs |> Enum.reject(fn({:ok, {status, _}}) -> status == :none end)
if is_boolean(expected_result) do
assert expected_result == !Enum.empty?(result)
else
assert expected_result == Enum.count(result)
end
end
test "enqueue/dequeue single queue" do
JobQueue.enqueue(:testredis, "test", "default", MyWorker, [], [])
[{:ok, {deq, _}}] = JobQueue.dequeue(:testredis, "test", @host, ["default"])
assert deq != :none
[{:ok, {deq, _}}] = JobQueue.dequeue(:testredis, "test", @host, ["default"])
assert deq == :none
end
test "enqueue/dequeue multi queue" do
JobQueue.enqueue(:testredis, "test", "default", MyWorker, [], [])
JobQueue.enqueue(:testredis, "test", "myqueue", MyWorker, [], [])
assert_dequeue_job(["default", "myqueue"], 2)
assert_dequeue_job(["default", "myqueue"], false)
end
test "backup queue" do
JobQueue.enqueue(:testredis, "test", "default", MyWorker, [], [])
JobQueue.enqueue(:testredis, "test", "default", MyWorker, [], [])
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], false)
JobQueue.re_enqueue_backup(:testredis, "test", @host, "default")
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], false)
end
test "remove from backup queue" do
JobQueue.enqueue(:testredis, "test", "default", MyWorker, [], [])
JobQueue.enqueue(:testredis, "test", "default", MyWorker, [], [])
[{:ok, {job, _}}] = JobQueue.dequeue(:testredis, "test", @host, ["default"])
assert_dequeue_job(["default"], true)
# remove job from queue
JobQueue.remove_job_from_backup(:testredis, "test", @host, "default", job)
# should only have 1 job now
JobQueue.re_enqueue_backup(:testredis, "test", @host, "default")
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], false)
end
test "scheduler_dequeue single queue" do
JobQueue.enqueue_in(:testredis, "test", "default", 0, MyWorker, [], [])
JobQueue.enqueue_in(:testredis, "test", "default", 0, MyWorker, [], [])
assert JobQueue.scheduler_dequeue(:testredis, "test") == 2
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], false)
end
test "scheduler_dequeue multi queue" do
JobQueue.enqueue_in(:testredis, "test", "default", -1, MyWorker, [], [])
JobQueue.enqueue_in(:testredis, "test", "myqueue", -1, MyWorker, [], [])
assert JobQueue.scheduler_dequeue(:testredis, "test") == 2
assert_dequeue_job(["default", "myqueue"], 2)
assert_dequeue_job(["default", "myqueue"], false)
end
test "scheduler_dequeue enqueue_at" do
JobQueue.enqueue_at(:testredis, "test", "default", DateTime.utc_now, MyWorker, [], [])
{jid, job_serialized} = JobQueue.to_job_serialized("retry", MyWorker, [], retry: true)
JobQueue.enqueue_job_at(:testredis, "test", job_serialized, jid, DateTime.utc_now, "test:retry")
assert JobQueue.scheduler_dequeue(:testredis, "test") == 2
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], false)
assert_dequeue_job(["retry"], true)
assert_dequeue_job(["retry"], false)
end
test "retry job" do
with_application_env(:exq, :max_retries, 1, fn ->
JobQueue.retry_or_fail_job(
:testredis,
"test",
%{
retry_count: 0,
retry: true,
queue: "default",
class: "MyWorker",
jid: UUID.uuid4,
error_class: nil,
error_message: "failed",
failed_at: Time.unix_seconds,
enqueued_at: Time.unix_seconds,
finished_at: nil,
processor: nil,
args: []
},
%RuntimeError{}
)
assert JobQueue.queue_size(:testredis, "test", :retry) == 1
end)
end
test "scheduler_dequeue max_score" do
add_usecs = fn(time, offset) ->
base = time |> DateTime.to_unix(:microseconds)
DateTime.from_unix!(base + offset, :microseconds)
end
JobQueue.enqueue_in(:testredis, "test", "default", 300, MyWorker, [], [])
now = DateTime.utc_now
time1 = add_usecs.(now, 140_000_000)
JobQueue.enqueue_at(:testredis, "test", "default", time1, MyWorker, [], [])
time2 = add_usecs.(now, 150_000_000)
JobQueue.enqueue_at(:testredis, "test", "default", time2, MyWorker, [], [])
time2a = add_usecs.(now, 151_000_000)
time2b = add_usecs.(now, 159_000_000)
time3 = add_usecs.(now, 160_000_000)
JobQueue.enqueue_at(:testredis, "test", "default", time3, MyWorker, [], [])
time4 = add_usecs.(now, 160_000_001)
JobQueue.enqueue_at(:testredis, "test", "default", time4, MyWorker, [], [])
time5 = add_usecs.(now, 300_000_000)
api_state = %Exq.Api.Server.State{redis: :testredis, namespace: "test"}
assert JobQueue.queue_size(api_state.redis, api_state.namespace, "default") == 0
assert JobQueue.queue_size(api_state.redis, api_state.namespace, :scheduled) == 5
assert JobQueue.scheduler_dequeue(:testredis, "test", Time.time_to_score(time2a)) == 2
assert JobQueue.scheduler_dequeue(:testredis, "test", Time.time_to_score(time2b)) == 0
assert JobQueue.scheduler_dequeue(:testredis, "test", Time.time_to_score(time3)) == 1
assert JobQueue.scheduler_dequeue(:testredis, "test", Time.time_to_score(time3)) == 0
assert JobQueue.scheduler_dequeue(:testredis, "test", Time.time_to_score(time4)) == 1
assert JobQueue.scheduler_dequeue(:testredis, "test", Time.time_to_score(time5)) == 1
assert JobQueue.queue_size(api_state.redis, api_state.namespace, "default") == 5
assert JobQueue.queue_size(api_state.redis, api_state.namespace, :scheduled) == 0
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], true)
assert_dequeue_job(["default"], false)
end
test "full_key" do
assert JobQueue.full_key("exq","k1") == "exq:k1"
assert JobQueue.full_key("","k1") == "k1"
assert JobQueue.full_key(nil,"k1") == "k1"
end
test "creates and returns a jid" do
{:ok, jid} = JobQueue.enqueue(:testredis, "test", "default", MyWorker, [], [])
assert jid != nil
[{:ok, {job_str, _}}] = JobQueue.dequeue(:testredis, "test", @host, ["default"])
job = Poison.decode!(job_str, as: %Exq.Support.Job{})
assert job.jid == jid
assert job.retry == Exq.Support.Config.get(:max_retries)
end
test "to_job_serialized using module atom" do
{_jid, serialized} = JobQueue.to_job_serialized("default", MyWorker, [], max_retries: 0)
job = Job.decode(serialized)
assert job.class == "MyWorker"
assert job.retry == 0
end
test "to_job_serialized using module string" do
{_jid, serialized} = JobQueue.to_job_serialized("default", "MyWorker/perform", [], max_retries: 10)
job = Job.decode(serialized)
assert job.class == "MyWorker/perform"
assert job.retry == 10
end
end