Packages
quantum
2.1.0
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: task supervisor could fail in cluster
Current section
Files
Jump to
Current section
Files
lib/quantum/execution_broadcaster.ex
defmodule Quantum.ExecutionBroadcaster do
@moduledoc """
Receives Added / Removed Jobs, Broadcasts Executions of Jobs
"""
use GenStage
require Logger
alias Quantum.{Job, Util, DateLibrary}
alias Crontab.{Scheduler, CronExpression}
@doc """
Start Stage
### Arguments
* `name` - The name of the stage
* `job_broadcaster` - The name of the stage to listen to
"""
@spec start_link(GenServer.server, GenServer.server) :: GenServer.on_start
def start_link(name, job_broadcaster) do
__MODULE__
|> GenStage.start_link(job_broadcaster, [name: name])
|> Util.start_or_link
end
@doc false
@spec child_spec({GenServer.server, GenServer.server}) :: Supervisor.child_spec
def child_spec({name, job_broadcaster}) do
%{super([]) | start: {__MODULE__, :start_link, [name, job_broadcaster]}}
end
@doc false
def init(job_broadcaster) do
state = %{jobs: [], time: NaiveDateTime.utc_now, timer: nil}
{:producer_consumer, state, subscribe_to: [job_broadcaster]}
end
def handle_events(events, _, state) do
reboot_add_events = events
|> Enum.filter(&add_reboot_event?/1)
|> Enum.map(fn {:add, job} -> {:execute, job} end)
state = events
|> Enum.reject(&add_reboot_event?/1)
|> Enum.reduce(state, &handle_event/2)
|> sort_state
|> reset_timer
{:noreply, reboot_add_events, state}
end
def handle_info(:execute, %{jobs: [{time_to_execute, jobs_to_execute} | tail]} = state) do
state = state
|> Map.put(:timer, nil)
|> Map.put(:jobs, tail)
|> Map.put(:time, NaiveDateTime.add(time_to_execute, 1, :second))
state = jobs_to_execute
|> Enum.reduce(state, &add_job_to_state/2)
|> sort_state
|> reset_timer
{:noreply, Enum.map(jobs_to_execute, fn job -> {:execute, job} end), state}
end
defp handle_event({:add, job}, state) do
add_job_to_state(job, state)
end
defp handle_event({:remove, name}, %{jobs: jobs} = state) do
jobs = jobs
|> Enum.map(fn {date, job_list} ->
{date, Enum.reject(job_list, &(&1.name == name))}
end)
|> Enum.reject(fn
{_, []} -> true
{_, _} -> false
end)
%{state | jobs: jobs}
end
defp add_job_to_state(%Job{schedule: schedule, timezone: timezone, name: name} = job, %{time: time} = state) do
case Scheduler.get_next_run_date(schedule, DateLibrary.to_tz!(time, timezone)) do
{:ok, date} ->
add_to_state(state, DateLibrary.to_utc!(date, timezone), job)
_ ->
Logger.warn """
Invalid Schedule #{inspect schedule} provided for job #{inspect name}.
No matching dates found. The job was removed.
"""
state
end
rescue
error ->
Logger.error("Invalid Timezone #{inspect timezone} provided for job #{inspect name}.",
job: job, error: error)
state
end
defp sort_state(%{jobs: jobs} = state) do
%{state | jobs: Enum.sort_by(jobs, fn {date, _} -> NaiveDateTime.to_erl(date) end)}
end
defp add_to_state(%{jobs: jobs} = state, date, job) do
%{state | jobs: case Enum.find_index(jobs, fn {run_date, _} -> run_date == date end) do
nil ->
[{date, [job]} | jobs]
index ->
List.update_at(jobs, index, fn {run_date, old} -> {run_date, [job | old]} end)
end}
end
defp reset_timer(%{timer: nil, jobs: []} = state) do
state
end
defp reset_timer(%{timer: {timer, _}, jobs: []} = state) do
Process.cancel_timer(timer)
Map.put(state, :timer, nil)
end
defp reset_timer(%{timer: nil, jobs: jobs} = state) do
run_date = next_run_date(jobs)
timer = case NaiveDateTime.compare(run_date, NaiveDateTime.utc_now) do
:gt ->
monotonic_time = run_date
|> DateTime.from_naive!("Etc/UTC")
|> DateTime.to_unix(:millisecond)
|> Kernel.-(System.time_offset(:millisecond))
Process.send_after(self(), :execute, monotonic_time, abs: true)
_ ->
send(self(), :execute)
nil
end
Map.put(state, :timer, {timer, run_date})
end
defp reset_timer(%{timer: {timer, old_date}, jobs: jobs} = state) do
run_date = next_run_date(jobs)
case NaiveDateTime.compare(run_date, old_date) do
:gt ->
Process.cancel_timer(timer)
reset_timer(Map.put(state, :timer, nil))
_ ->
state
end
end
defp next_run_date(jobs) do
[{run_date, _} | _rest] = jobs
run_date
end
defp add_reboot_event?({:add, %Job{schedule: %CronExpression{reboot: true}}}), do: true
defp add_reboot_event?(_), do: false
end