Packages
exq
0.1.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/redis_queue.ex
defmodule Exq.RedisQueue do
use Timex
@default_queue "default"
def find_job(redis, namespace, jid, queue) do
jobs = Exq.Redis.lrange!(redis, queue_key(namespace, queue))
finder = fn({j, idx}) ->
job = Exq.Job.from_json(j)
job.jid == jid
end
error = Enum.find(Enum.with_index(jobs), finder)
case error do
nil ->
{:not_found, nil}
_ ->
{job, idx} = error
{:ok, job, idx}
end
end
def enqueue(redis, namespace, queue, worker, args) do
{jid, job} = job_json(queue, worker, args)
[{:ok, _}, {:ok, _}] = :eredis.qp(redis, [
["SADD", full_key(namespace, "queues"), queue],
["RPUSH", queue_key(namespace, queue), job]])
jid
end
def dequeue(redis, namespace, queues) when is_list(queues) do
dequeue_random(redis, namespace, queues)
end
def dequeue(redis, namespace, queue) do
Exq.Redis.lpop!(redis, queue_key(namespace, queue))
end
def full_key(namespace, key) do
"#{namespace}:#{key}"
end
def queue_key(namespace, queue) do
full_key(namespace, "queue:#{queue}")
end
defp dequeue_random(redis, namespace, []) do
nil
end
defp dequeue_random(redis, namespace, queues) do
[h | rq] = Exq.Shuffle.shuffle(queues)
case dequeue(redis, namespace, h) do
nil -> dequeue_random(redis, namespace, rq)
job -> job
end
end
defp job_json(queue, worker, args) do
jid = UUID.uuid4
job = Enum.into([{:queue, queue}, {:class, worker}, {:args, args}, {:jid, jid}, {:enqueued_at, DateFormat.format!(Date.local, "{ISO}")}], HashDict.new)
{jid, Exq.Json.encode(job)}
end
end