Packages
exq
0.13.2
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/scheduler/server.ex
defmodule Exq.Scheduler.Server do
@moduledoc """
The Scheduler is responsible for monitoring the `schedule` and `retry` queues.
These queues use a Redis sorted set (term?) to schedule and pick off due jobs.
Once a job is at or past it's execution date, the Scheduler moves the job into the
live execution queue.
Runs on a timed loop according to `scheduler_poll_timeout`.
## Initialization:
* `:name` - Name of target registered process
* `:namespace` - Redis namespace to store all data under. Defaults to "exq".
* `:queues` - Array of currently active queues (TODO: Remove, I suspect it's not needed).
* `:redis` - pid of Redis process.
* `:scheduler_poll_timeout` - How often to poll Redis for scheduled / retry jobs.
"""
require Logger
use GenServer
defmodule State do
defstruct redis: nil, namespace: nil, queues: nil, scheduler_poll_timeout: nil
end
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: server_name(opts[:name]))
end
def start_timeout(pid) do
GenServer.cast(pid, :start_timeout)
end
def server_name(name) do
name = name || Exq.Support.Config.get(:name)
"#{name}.Scheduler" |> String.to_atom()
end
## ===========================================================
## gen server callbacks
## ===========================================================
def init(opts) do
state = %State{
redis: opts[:redis],
namespace: opts[:namespace],
queues: opts[:queues],
scheduler_poll_timeout: opts[:scheduler_poll_timeout]
}
start_timeout(self())
{:ok, state}
end
def handle_cast(:start_timeout, state) do
handle_info(:timeout, state)
end
def handle_info(:timeout, state) do
{updated_state, timeout} = dequeue(state)
{:noreply, updated_state, timeout}
end
## ===========================================================
## Internal Functions
## ===========================================================
@doc """
Dequeue any active jobs in the scheduled and retry queues, and enqueue them to live queue.
"""
def dequeue(state) do
Exq.Redis.JobQueue.scheduler_dequeue(state.redis, state.namespace)
{state, state.scheduler_poll_timeout}
end
end