Packages
oban
1.0.0
2.23.0
2.22.1
2.22.0
2.21.1
2.21.0
2.20.3
2.20.2
2.20.1
2.20.0
2.19.4
2.19.3
2.19.2
2.19.1
2.19.0
2.18.3
2.18.2
2.18.1
2.18.0
2.17.12
2.17.11
2.17.10
2.17.9
2.17.8
2.17.7
2.17.6
2.17.5
2.17.4
2.17.3
2.17.2
2.17.1
2.17.0
2.16.3
2.16.2
2.16.1
2.16.0
2.15.4
2.15.3
2.15.2
2.15.1
2.15.0
2.14.2
2.14.1
2.14.0
2.13.6
2.13.5
2.13.4
2.13.3
2.13.2
2.13.1
2.13.0
2.12.1
2.12.0
2.11.3
2.11.2
2.11.1
2.11.0
2.10.1
2.10.0
retired
2.9.2
2.9.1
2.9.0
2.8.0
2.7.2
2.7.1
2.7.0
2.6.1
2.6.0
2.5.0
2.4.3
2.4.2
2.4.1
2.4.0
2.3.4
2.3.3
2.3.2
2.3.1
2.3.0
2.2.0
2.1.0
2.0.0
2.0.0-rc.3
2.0.0-rc.2
2.0.0-rc.1
2.0.0-rc.0
1.2.0
1.1.0
1.0.0
1.0.0-rc.2
1.0.0-rc.1
0.12.1
0.12.0
0.11.1
0.11.0
0.10.1
0.10.0
0.9.0
0.8.1
0.8.0
0.7.1
0.7.0
0.6.0
0.5.0
0.4.0
0.3.0
0.2.0
0.1.0
Robust job processing, backed by modern PostgreSQL, SQLite3, and MySQL.
Current section
Files
Jump to
Current section
Files
lib/oban/crontab/scheduler.ex
defmodule Oban.Crontab.Scheduler do
@moduledoc false
use GenServer
import Oban.Breaker, only: [open_circuit: 1, trip_errors: 0, trip_circuit: 2]
alias Oban.Crontab.Cron
alias Oban.{Config, Query, Worker}
# Make each job unique for 59 seconds to prevent double-enqueue if the node or scheduler
# crashes. The minimum resolution for our cron jobs is 1 minute, so there is potentially
# a one second window where a double enqueue can happen.
@unique [period: 59]
# Use an extremely large 64bit integer as the key value to prevent any chance of intersecting
# with application keyspace.
@lock_key 1_149_979_440_242_868_330
@type option :: {:name, module()} | {:conf, Config.t()}
defmodule State do
@moduledoc false
@enforce_keys [:conf]
defstruct [
:conf,
:name,
:poll_ref,
circuit: :enabled,
poll_interval: :timer.seconds(60)
]
end
@spec start_link([option()]) :: GenServer.on_start()
def start_link(opts) when is_list(opts) do
name = Keyword.get(opts, :name, __MODULE__)
GenServer.start_link(__MODULE__, opts[:conf], name: name)
end
@impl GenServer
def init(%Config{crontab: []}) do
:ignore
end
def init(%Config{} = conf) do
Process.flag(:trap_exit, true)
{:ok, struct!(State, conf: conf), {:continue, :start}}
end
@impl GenServer
def handle_continue(:start, state) do
handle_info(:poll, state)
end
@impl GenServer
def terminate(_reason, %State{poll_ref: poll_ref}) do
if not is_nil(poll_ref), do: Process.cancel_timer(poll_ref)
:ok
end
@impl GenServer
def handle_info(:poll, state) do
state =
state
|> send_poll_after()
|> lock_and_enqueue()
{:noreply, state}
end
def handle_info({:EXIT, _pid, error}, %State{} = state) do
{:noreply, trip_circuit(error, state)}
end
def handle_info(:reset_circuit, state) do
{:noreply, open_circuit(state)}
end
defp send_poll_after(%State{poll_interval: interval} = state) do
ref = Process.send_after(self(), :poll, interval)
%{state | poll_ref: ref}
end
defp lock_and_enqueue(%State{circuit: :disabled} = state), do: state
defp lock_and_enqueue(%State{conf: conf, poll_interval: timeout} = state) do
%Config{repo: repo, verbose: verbose} = conf
repo.transaction(
fn -> if acquire_lock?(conf), do: enqueue_jobs(conf) end,
log: verbose,
timeout: timeout
)
state
rescue
exception in trip_errors() -> trip_circuit(exception, state)
end
defp acquire_lock?(%Config{repo: repo, verbose: verbose}) do
%{rows: [[locked?]]} =
repo.query!(
"SELECT pg_try_advisory_xact_lock($1)",
[@lock_key],
log: verbose
)
locked?
end
defp enqueue_jobs(%Config{crontab: crontab, timezone: timezone} = conf) do
{:ok, datetime} = DateTime.now(timezone)
for {cron, worker, opts} <- crontab, Cron.now?(cron, datetime) do
{args, opts} = Keyword.pop(opts, :args, %{})
opts = unique_opts(worker.__opts__(), opts)
{:ok, _job} = Query.fetch_or_insert_job(conf, worker.new(args, opts))
end
end
# Ensure that `unique` is deep merged and the default period has the lowest priority.
defp unique_opts(worker_opts, crontab_opts) do
[unique: @unique]
|> Keyword.merge(worker_opts, &Worker.resolve_opts/3)
|> Keyword.merge(crontab_opts, &Worker.resolve_opts/3)
end
end