Packages
quantum
2.3.1
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.
Retired package: Release invalid
Current section
Files
Jump to
Current section
Files
lib/quantum/executor_supervisor.ex
defmodule Quantum.ExecutorSupervisor do
@moduledoc """
This `ConsumerSupervisor` is responsible to start a job for every execute event.
"""
use ConsumerSupervisor
alias Quantum.Util
@spec start_link(
GenServer.server(),
GenServer.server(),
GenServer.server(),
GenServer.server(),
boolean()
) :: GenServer.on_start()
def start_link(name, execution_broadcaster, task_supervisor, task_registry, debug_logging) do
ConsumerSupervisor.start_link(
__MODULE__,
%{
execution_broadcaster: execution_broadcaster,
task_supervisor: task_supervisor,
task_registry: task_registry,
cluster_task_supervisor_registry: nil,
debug_logging: debug_logging
},
name: name
)
end
@spec start_link(%{
name: GenServer.server(),
execution_broadcaster: GenServer.server(),
task_supervisor: GenServer.server(),
task_registry: GenServer.server(),
cluster_task_supervisor_registry: nil | GenServer.server(),
debug_logging: boolean()
}) :: GenServer.on_start()
def start_link(opts) do
ConsumerSupervisor.start_link(
__MODULE__,
Map.take(opts, [
:execution_broadcaster,
:task_supervisor,
:task_registry,
:cluster_task_supervisor_registry,
:debug_logging
]),
name: Map.fetch!(opts, :name)
)
end
# credo:disable-for-next-line Credo.Check.Design.TagTODO
# TODO: Remove when gen_stage:0.12 support is dropped
if Util.gen_stage_v12?() do
def init(
%{
execution_broadcaster: execution_broadcaster
} = opts
) do
ConsumerSupervisor.init(
{Quantum.Executor,
Map.take(opts, [
:task_supervisor,
:task_registry,
:debug_logging,
:cluster_task_supervisor_registry
])},
strategy: :one_for_one,
subscribe_to: [{execution_broadcaster, max_demand: 50}]
)
end
else
def init(
%{
execution_broadcaster: execution_broadcaster
} = opts
) do
ConsumerSupervisor.init(
[
{Quantum.Executor,
Map.take(opts, [
:task_supervisor,
:task_registry,
:debug_logging,
:cluster_task_supervisor_registry
])}
],
strategy: :one_for_one,
subscribe_to: [{execution_broadcaster, max_demand: 50}]
)
end
end
@doc false
@spec child_spec(
{
GenServer.server(),
GenServer.server(),
GenServer.server(),
GenServer.server(),
boolean()
}
| {
GenServer.server(),
GenServer.server(),
GenServer.server(),
GenServer.server(),
GenServer.server(),
boolean()
}
) :: Supervisor.child_spec()
def child_spec({name, execution_broadcaster, task_supervisor, task_registry, debug_logging}),
do:
child_spec(
{name, execution_broadcaster, task_supervisor, task_registry, nil, debug_logging}
)
def child_spec(
{name, execution_broadcaster, task_supervisor, task_registry,
cluster_task_supervisor_registry, debug_logging}
) do
%{
super([])
| start: {
__MODULE__,
:start_link,
[
%{
name: name,
execution_broadcaster: execution_broadcaster,
task_supervisor: task_supervisor,
task_registry: task_registry,
cluster_task_supervisor_registry: cluster_task_supervisor_registry,
debug_logging: debug_logging
}
]
}
}
end
end