Packages
quantum
3.0.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/executor.ex
defmodule Quantum.Executor do
@moduledoc false
# Task to actually execute a Task
use Task
require Logger
alias Quantum.{
Job,
NodeSelectorBroadcaster.Event,
TaskRegistry
}
alias __MODULE__.StartOpts
@spec start_link(StartOpts.t(), Event.t()) :: {:ok, pid}
def start_link(opts, %Event{job: job, node: node}) do
Task.start_link(fn ->
execute(opts, job, node)
end)
end
@spec execute(StartOpts.t(), Job.t(), Node.t()) :: :ok
# Execute task on all given nodes without checking for overlap
defp execute(
%StartOpts{task_supervisor_reference: task_supervisor, debug_logging: debug_logging},
%Job{overlap: true} = job,
node
) do
run(node, job, task_supervisor, debug_logging)
:ok
end
# Execute task on all given nodes with checking for overlap
defp execute(
%StartOpts{
task_supervisor_reference: task_supervisor,
task_registry_reference: task_registry,
debug_logging: debug_logging
},
%Job{overlap: false, name: job_name} = job,
node
) do
debug_logging &&
Logger.debug(fn ->
"[#{inspect(Node.self())}][#{__MODULE__}] Start execution of job #{inspect(job_name)}"
end)
case TaskRegistry.mark_running(task_registry, job_name, node) do
:marked_running ->
%Task{ref: ref} = run(node, job, task_supervisor, debug_logging)
receive do
{^ref, _} ->
TaskRegistry.mark_finished(task_registry, job_name, node)
{:DOWN, ^ref, _, _, _} ->
TaskRegistry.mark_finished(task_registry, job_name, node)
:ok
end
_ ->
:ok
end
end
# Ececute the given function on a given node via the task supervisor
@spec run(Node.t(), Job.t(), GenServer.server(), boolean()) :: Task.t()
defp run(node, %{name: job_name, task: task}, task_supervisor, debug_logging) do
debug_logging &&
Logger.debug(fn ->
"[#{inspect(Node.self())}][#{__MODULE__}] Task for job #{inspect(job_name)} started on node #{
inspect(node)
}"
end)
Task.Supervisor.async_nolink({task_supervisor, node}, fn ->
debug_logging &&
Logger.debug(fn ->
"[#{inspect(Node.self())}][#{__MODULE__}] Execute started for job #{inspect(job_name)}"
end)
try do
execute_task(task)
catch
type, value ->
debug_logging &&
Logger.debug(fn ->
"[#{inspect(Node.self())}][#{__MODULE__}] Execution ended for job #{
inspect(job_name)
}, which failed due to: #{Exception.format(type, value, __STACKTRACE__)}"
end)
else
result ->
debug_logging &&
Logger.debug(fn ->
"[#{inspect(Node.self())}][#{__MODULE__}] Execution ended for job #{
inspect(job_name)
}, which yielded result: #{inspect(result)}"
end)
end
:ok
end)
end
# Run function
@spec execute_task(Quantum.Job.task()) :: any
defp execute_task({mod, fun, args}) do
:erlang.apply(mod, fun, args)
end
defp execute_task(fun) when is_function(fun, 0) do
fun.()
end
end