Packages
quantum
2.3.2
3.5.3
3.5.2
retired
3.5.1
retired
3.5.0
3.4.0
3.3.0
3.2.0
3.1.0
3.0.2
3.0.1
3.0.0
3.0.0-rc.3
3.0.0-rc.2
3.0.0-rc.1
2.4.0
2.3.4
2.3.3
2.3.2
2.3.1
retired
2.3.0
retired
2.2.7
2.2.6
2.2.5
retired
2.2.4
retired
2.2.3
retired
2.2.2
retired
2.2.1
retired
2.2.0
retired
2.1.3
retired
2.1.2
retired
2.1.1
retired
2.1.0
retired
2.1.0-beta.1
retired
2.0.4
retired
2.0.3
retired
2.0.2
retired
2.0.1
retired
2.0.0
retired
2.0.0-beta.2
retired
2.0.0-beta.1
retired
1.9.3
retired
1.9.2
retired
1.9.1
retired
1.9.0
retired
1.8.1
retired
1.8.0
retired
1.7.1
retired
1.7.0
retired
1.6.1
retired
1.6.0
retired
1.5.0
retired
1.4.0
retired
1.3.2
retired
1.3.1
retired
1.3.0
retired
1.2.4
retired
1.2.3
retired
1.2.2
retired
1.2.1
retired
1.2.0
retired
1.1.0
retired
1.0.4
retired
1.0.3
retired
1.0.2
retired
1.0.1
retired
1.0.0
retired
Cron-like job scheduler for Elixir.
Current section
Files
Jump to
Current section
Files
lib/quantum.ex
defmodule Quantum do
@moduledoc """
Contains config functions to aid the rest of the lib
"""
require Logger
alias Quantum.{Job, Normalizer, RunStrategy.Random, Storage.Noop}
@defaults [
global: false,
cron: [],
timeout: 5_000,
schedule: nil,
overlap: true,
timezone: :utc,
run_strategy: {Random, :cluster},
debug_logging: true
]
@doc """
Retrieves only scheduler related configuration.
"""
def scheduler_config(quantum, otp_app, custom) do
config =
@defaults
|> Keyword.merge(Application.get_env(otp_app, quantum, []))
|> Keyword.merge(custom)
|> Keyword.merge(otp_app: otp_app, quantum: quantum)
global = Keyword.fetch!(config, :global)
job_broadcaster_name = Module.concat(quantum, JobBroadcaster)
job_broadcaster_reference =
if global,
do: {:via, :swarm, job_broadcaster_name},
else: job_broadcaster_name
execution_broadcaster_name = Module.concat(quantum, ExecutionBroadcaster)
execution_broadcaster_reference =
if global,
do: {:via, :swarm, execution_broadcaster_name},
else: execution_broadcaster_name
executor_supervisor_name = Module.concat(quantum, ExecutorSupervisor)
executor_supervisor_reference = executor_supervisor_name
task_registry_name = Module.concat(quantum, TaskRegistry)
task_registry_reference =
if global,
do: {:via, :swarm, task_registry_name},
else: task_registry_name
cluster_task_supervisor_registry_name = Module.concat(quantum, ClusterTaskSupervisorRegistry)
cluster_task_supervisor_registry_reference = cluster_task_supervisor_registry_name
task_supervisor_name = Module.concat(quantum, Task.Supervisor)
task_supervisor_reference = task_supervisor_name
config
|> Keyword.put_new(:quantum, quantum)
|> Keyword.put_new(:scheduler, quantum)
|> update_in([:schedule], &Normalizer.normalize_schedule/1)
|> Keyword.put_new(:job_broadcaster_name, job_broadcaster_name)
|> Keyword.put_new(:job_broadcaster_reference, job_broadcaster_reference)
|> Keyword.put_new(:execution_broadcaster_name, execution_broadcaster_name)
|> Keyword.put_new(:execution_broadcaster_reference, execution_broadcaster_reference)
|> Keyword.put_new(:executor_supervisor_name, executor_supervisor_name)
|> Keyword.put_new(:executor_supervisor_reference, executor_supervisor_reference)
|> Keyword.put_new(:task_registry_name, task_registry_name)
|> Keyword.put_new(:task_registry_reference, task_registry_reference)
|> Keyword.put_new(:task_supervisor_name, task_supervisor_name)
|> Keyword.put_new(:task_supervisor_reference, task_supervisor_reference)
|> Keyword.put_new(
:cluster_task_supervisor_registry_name,
cluster_task_supervisor_registry_name
)
|> Keyword.put_new(
:cluster_task_supervisor_registry_reference,
cluster_task_supervisor_registry_reference
)
|> Keyword.put_new(:storage, Noop)
end
@doc """
Retrieves the comprehensive runtime configuration.
"""
def runtime_config(quantum, otp_app, custom) do
config = scheduler_config(quantum, otp_app, custom)
# Load Jobs from Config
jobs =
config
|> Keyword.get(:jobs, [])
|> Enum.map(&Normalizer.normalize(quantum.new_job(config), &1))
|> remove_jobs_with_duplicate_names(quantum)
Keyword.put(config, :jobs, jobs)
end
defp remove_jobs_with_duplicate_names(job_list, quantum) do
job_list
|> Enum.reduce(%{}, fn %Job{name: name} = job, acc ->
if Enum.member?(Map.keys(acc), name) do
Logger.warn(
"Job with name '#{name}' of quantum '#{quantum}' not started due to duplicate job name"
)
acc
else
Map.put_new(acc, name, job)
end
end)
|> Map.values()
end
end