Packages
exq
0.10.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/support/opts.ex
defmodule Exq.Support.Opts do
alias Exq.Support.Coercion
alias Exq.Support.Config
@doc """
Return top supervisor's name default is Exq.Sup
"""
def top_supervisor(name) do
name = name || Config.get(:name)
"#{name}.Sup" |> String.to_atom
end
defp conform_opts(opts) do
mode = opts[:mode] || Config.get(:mode)
redis = redis_client_name(opts[:name])
opts = [{:redis, redis}|opts]
redis_opts = redis_opts(opts)
connection_opts = connection_opts(opts)
server_opts = server_opts(mode, opts)
{redis_opts, connection_opts, server_opts}
end
def redis_client_name(name) do
name = name || Config.get(:name)
"#{name}.Redis.Client" |> String.to_atom
end
def redis_opts(opts \\ []) do
if url = opts[:url] || Config.get(:url) do
url
else
host = opts[:host] || Config.get(:host)
port = Coercion.to_integer(opts[:port] || Config.get(:port))
database = Coercion.to_integer(opts[:database] || Config.get(:database))
password = opts[:password] || Config.get(:password)
[host: host, port: port, database: database, password: password]
end
end
@doc """
Return {redis_module, redis_args, gen_server_opts}
"""
def redis_worker_opts(opts) do
{redis_opts, connection_opts, opts} = conform_opts(opts)
case Config.get(:redis_worker) do
{module, args} -> {module, args, opts}
_ -> {Redix, [redis_opts, connection_opts], opts}
end
end
def redis_worker_module() do
case Config.get(:redis_worker) do
{module, _args} -> module
_ -> Redix
end
end
def connection_opts(opts \\ []) do
reconnect_on_sleep = opts[:reconnect_on_sleep] || Config.get(:reconnect_on_sleep)
timeout = opts[:redis_timeout] || Config.get(:redis_timeout)
socket_opts = opts[:socket_opts] || Config.get(:socket_opts) || []
[backoff: reconnect_on_sleep, timeout: timeout, name: opts[:redis], socket_opts: socket_opts]
end
defp server_opts(:default, opts) do
scheduler_enable = opts[:scheduler_enable] || Config.get(:scheduler_enable)
namespace = opts[:namespace] || Config.get(:namespace)
scheduler_poll_timeout = opts[:scheduler_poll_timeout] || Config.get(:scheduler_poll_timeout)
poll_timeout = opts[:poll_timeout] || Config.get(:poll_timeout)
shutdown_timeout = opts[:shutdown_timeout] || Config.get(:shutdown_timeout)
enqueuer = Exq.Enqueuer.Server.server_name(opts[:name])
stats = Exq.Stats.Server.server_name(opts[:name])
scheduler = Exq.Scheduler.Server.server_name(opts[:name])
workers_sup = Exq.Worker.Supervisor.supervisor_name(opts[:name])
middleware = Exq.Middleware.Server.server_name(opts[:name])
metadata = Exq.Worker.Metadata.server_name(opts[:name])
queue_configs = opts[:queues] || Config.get(:queues)
per_queue_concurrency = opts[:concurrency] || Config.get(:concurrency)
queues = get_queues(queue_configs)
concurrency = get_concurrency(queue_configs, per_queue_concurrency)
default_middleware = Config.get(:middleware)
[scheduler_enable: scheduler_enable, namespace: namespace,
scheduler_poll_timeout: scheduler_poll_timeout,workers_sup: workers_sup,
poll_timeout: poll_timeout, enqueuer: enqueuer, metadata: metadata,
stats: stats, name: opts[:name], scheduler: scheduler, queues:
queues, redis: opts[:redis], concurrency: concurrency,
middleware: middleware, default_middleware: default_middleware,
mode: :default, shutdown_timeout: shutdown_timeout]
end
defp server_opts(mode, opts) do
namespace = opts[:namespace] || Config.get(:namespace)
[name: opts[:name], namespace: namespace, redis: opts[:redis], mode: mode]
end
defp get_queues(queue_configs) do
Enum.map(queue_configs, fn queue_config ->
case queue_config do
{queue, _concurrency} -> queue
queue -> queue
end
end)
end
defp get_concurrency(queue_configs, per_queue_concurrency) do
Enum.map(queue_configs, fn (queue_config) ->
case queue_config do
{queue, concurrency} -> {queue, concurrency, 0}
queue -> {queue, per_queue_concurrency, 0}
end
end)
end
end