Packages
exq
0.6.1
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.ex
defmodule Exq do
require Logger
alias Exq.Support.Config
use Application
# OTP Application
def start(_type, _args) do
Exq.Manager.Supervisor.start_link
end
# Exq methods
def start(opts \\ []) do
Exq.Manager.Supervisor.start_link(opts)
end
def start_link(opts \\ []) do
Exq.Manager.Supervisor.start_link(opts)
end
def stop(pid) when is_pid(pid) do
Process.exit(pid, :shutdown)
end
def stop(name) when is_atom(name) do
case Process.whereis(Exq.Manager.Supervisor.supervisor_name(name)) do
nil -> :ok
pid ->
stop(pid)
end
end
@doc """
Enqueue a job immediately.
Expected args:
* `pid` - PID for Exq Manager or Enqueuer to handle this
* `queue` - Name of queue to use
* `worker` - Worker module to target
* `args` - Array of args to send to worker
Returns:
* `{:ok, jid}` if the job was enqueued successfully, with `jid` = Job ID.
* `{:error, reason}` if there was an error enqueueing job
"""
def enqueue(pid, queue, worker, args) do
GenServer.call(pid, {:enqueue, queue, worker, args}, Config.get(:genserver_timeout, 5000))
end
@doc """
Schedule a job to be enqueued at a specific time in the future.
Expected args:
* `pid` - PID for Exq Manager or Enqueuer to handle this
* `queue` - name of queue to use
* `time` - Time to enqueue
* `worker` - Worker module to target
* `args` - Array of args to send to worker
"""
def enqueue_at(pid, queue, time, worker, args) do
GenServer.call(pid, {:enqueue_at, queue, time, worker, args}, Config.get(:genserver_timeout, 5000))
end
@doc """
Schedule a job to be enqueued at in the future given by offset in milliseconds.
Expected args:
* `pid` - PID for Exq Manager or Enqueuer to handle this
* `queue` - Name of queue to use
* offset - Offset in milliseconds in the future to enqueue
* `worker` - Worker module to target
* `args` - Array of args to send to worker
"""
def enqueue_in(pid, queue, offset, worker, args) do
GenServer.call(pid, {:enqueue_in, queue, offset, worker, args}, Config.get(:genserver_timeout, 5000))
end
@doc """
Subscribe to a queue - ie. listen to queue for jobs
* `pid` - PID for Exq Manager or Enqueuer to handle this
* `queue` - Name of queue
* `concurrency` - Optional argument specifying max concurrency for queue
"""
def subscribe(pid, queue) do
GenServer.call(pid, {:subscribe, queue})
end
def subscribe(pid, queue, concurrency) do
GenServer.call(pid, {:subscribe, queue, concurrency})
end
@doc """
Unsubscribe from a queue - ie. stop listening to queue for jobs
* `pid` - PID for Exq Manager or Enqueuer to handle this
* `queue` - Name of queue
"""
def unsubscribe(pid, queue) do
GenServer.call(pid, {:unsubscribe, queue})
end
end