Packages
quantum
2.2.6
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/job_broadcaster.ex
defmodule Quantum.JobBroadcaster do
@moduledoc """
This Module is here to broadcast added / removed tabs into the execution pipeline.
"""
use GenStage
require Logger
alias Quantum.{Job, Util}
@doc """
Start Job Broadcaster
### Arguments
* `name` - Name of the GenStage
* `jobs` - Array of `Quantum.Job`
"""
@spec start_link(GenServer.server(), [Job.t()]) :: GenServer.on_start()
def start_link(name, jobs) do
__MODULE__
|> GenStage.start_link(jobs, name: name)
|> Util.start_or_link()
end
@doc false
@spec child_spec({GenServer.server(), [Job.t()]}) :: Supervisor.child_spec()
def child_spec({name, jobs}) do
%{super([]) | start: {__MODULE__, :start_link, [name, jobs]}}
end
@doc false
def init(jobs) do
state = %{
jobs: Enum.into(jobs, %{}, fn %{name: name} = job -> {name, job} end),
buffer: for(%{state: :active} = job <- jobs, do: {:add, job})
}
{:producer, state}
end
def handle_demand(demand, %{buffer: buffer} = state) do
{to_send, remaining} = Enum.split(buffer, demand)
{:noreply, to_send, %{state | buffer: remaining}}
end
def handle_cast({:add, %Job{state: :active, name: job_name} = job}, %{jobs: jobs} = state) do
Logger.debug(fn ->
"[#{inspect(Node.self())}][#{__MODULE__}] Adding job #{inspect(job_name)}"
end)
{:noreply, [{:add, job}], %{state | jobs: Map.put(jobs, job_name, job)}}
end
def handle_cast({:add, %Job{state: :inactive, name: job_name} = job}, %{jobs: jobs} = state) do
Logger.debug(fn ->
"[#{inspect(Node.self())}][#{__MODULE__}] Adding job #{inspect(job_name)}"
end)
{:noreply, [], %{state | jobs: Map.put(jobs, job_name, job)}}
end
def handle_cast({:delete, name}, %{jobs: jobs} = state) do
Logger.debug(fn ->
"[#{inspect(Node.self())}][#{__MODULE__}] Deleting job #{inspect(name)}"
end)
case Map.fetch(jobs, name) do
{:ok, %{state: :active}} ->
{:noreply, [{:remove, name}], %{state | jobs: Map.delete(jobs, name)}}
{:ok, %{state: :inactive}} ->
{:noreply, [], %{state | jobs: Map.delete(jobs, name)}}
:error ->
{:noreply, [], state}
end
end
def handle_cast({:change_state, name, new_state}, %{jobs: jobs} = state) do
Logger.debug(fn ->
"[#{inspect(Node.self())}][#{__MODULE__}] Change job state #{inspect(name)}"
end)
case Map.fetch(jobs, name) do
:error ->
{:noreply, [], state}
{:ok, %{state: ^new_state}} ->
{:noreply, [], state}
{:ok, job} ->
jobs = Map.update!(jobs, name, &Job.set_state(&1, new_state))
case new_state do
:active ->
{:noreply, [{:add, %{job | state: new_state}}], %{state | jobs: jobs}}
:inactive ->
{:noreply, [{:remove, name}], %{state | jobs: jobs}}
end
end
end
def handle_cast(:delete_all, %{jobs: jobs} = state) do
Logger.debug(fn ->
"[#{inspect(Node.self())}][#{__MODULE__}] Deleting all jobs"
end)
messages = for {name, %Job{state: :active}} <- jobs, do: {:remove, name}
{:noreply, messages, %{state | jobs: %{}}}
end
def handle_call(:jobs, _, %{jobs: jobs} = state), do: {:reply, Map.to_list(jobs), [], state}
def handle_call({:find_job, name}, _, %{jobs: jobs} = state),
do: {:reply, Map.get(jobs, name), [], state}
end