Packages
exq
0.14.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/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)
if is_binary(redis_opts) do
{Redix, [redis_opts, connection_opts], opts}
else
{Redix, [Keyword.merge(redis_opts, connection_opts)], opts}
end
end
def connection_opts(opts \\ []) do
redis_options = opts[:redis_options] || Config.get(:redis_options)
socket_opts = opts[:socket_opts] || Config.get(:socket_opts) || []
Keyword.merge(
[name: opts[:redis], socket_opts: socket_opts],
redis_options
)
end
defp server_opts(:default, opts) do
scheduler_enable =
Coercion.to_boolean(opts[:scheduler_enable] || Config.get(:scheduler_enable))
namespace = opts[:namespace] || Config.get(:namespace)
scheduler_poll_timeout =
Coercion.to_integer(opts[:scheduler_poll_timeout] || Config.get(:scheduler_poll_timeout))
poll_timeout = Coercion.to_integer(opts[:poll_timeout] || Config.get(:poll_timeout))
shutdown_timeout =
Coercion.to_integer(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] || get_config_concurrency()
queues = get_queues(queue_configs)
concurrency = get_concurrency(queue_configs, per_queue_concurrency)
default_middleware = Config.get(:middleware)
heartbeat_enable =
Coercion.to_boolean(Keyword.get(opts, :heartbeat_enable, Config.get(:heartbeat_enable)))
heartbeat_interval =
Coercion.to_integer(opts[:heartbeat_interval] || Config.get(:heartbeat_interval))
missed_heartbeats_allowed =
Coercion.to_integer(
opts[:missed_heartbeats_allowed] || Config.get(:missed_heartbeats_allowed)
)
[
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,
heartbeat_enable: heartbeat_enable,
heartbeat_interval: heartbeat_interval,
missed_heartbeats_allowed: missed_heartbeats_allowed
]
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_config_concurrency() do
cast_concurrency(Config.get(:concurrency))
end
defp get_concurrency(queue_configs, per_queue_concurrency) do
Enum.map(queue_configs, fn queue_config ->
case queue_config do
{queue, concurrency} -> {queue, cast_concurrency(concurrency), 0}
queue -> {queue, per_queue_concurrency, 0}
end
end)
end
defp cast_concurrency(:infinity), do: :infinity
defp cast_concurrency(:infinite), do: :infinity
defp cast_concurrency(x) when is_integer(x), do: x
defp cast_concurrency(x) when is_binary(x) do
case x |> String.trim() |> String.downcase() do
"infinity" -> :infinity
"infinite" -> :infinity
x -> Coercion.to_integer(x)
end
end
end