Packages
oban
2.20.2
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/stager.ex
defmodule Oban.Stager do
@moduledoc false
use GenServer
alias Oban.{Engine, Job, Notifier, Peer, Plugin, Registry, Repo}
alias __MODULE__, as: State
require Logger
@type option :: Plugin.option() | {:interval, pos_integer()}
defstruct [
:conf,
:timer,
interval: :timer.seconds(1),
limit: 5_000,
mode: :global
]
@spec start_link([option()]) :: GenServer.on_start()
def start_link(opts) do
{name, opts} = Keyword.pop(opts, :name)
conf = Keyword.fetch!(opts, :conf)
if conf.stage_interval == :infinity do
:ignore
else
state = %State{conf: conf, interval: conf.stage_interval}
GenServer.start_link(__MODULE__, state, name: name)
end
end
@impl GenServer
def init(state) do
Process.flag(:trap_exit, true)
# Init event is essential for auto-allow and backward compatibility.
:telemetry.execute([:oban, :plugin, :init], %{}, %{conf: state.conf, plugin: __MODULE__})
{:ok, schedule_staging(state)}
end
@impl GenServer
def terminate(_reason, %State{timer: timer}) do
if is_reference(timer), do: Process.cancel_timer(timer)
:ok
end
@impl GenServer
def handle_info(:stage, %State{} = state) do
state = check_mode(state)
meta = %{conf: state.conf, leader: Peer.leader?(state.conf), plugin: __MODULE__}
:telemetry.span([:oban, :plugin], meta, fn ->
case stage_and_notify(meta.leader, state) do
{:ok, staged} ->
{:ok, Map.merge(meta, %{staged_count: length(staged), staged_jobs: staged})}
{:error, error} ->
{:error, Map.put(meta, :error, error)}
end
end)
{:noreply, schedule_staging(state)}
end
def handle_info(message, state) do
Logger.warning(
[
message: "Received unexpected message: #{inspect(message)}",
source: :oban,
module: __MODULE__
],
domain: [:oban]
)
{:noreply, state}
end
defp stage_and_notify(true = _leader, state) do
Repo.transaction(state.conf, fn ->
{:ok, staged} = Engine.stage_jobs(state.conf, Job, limit: state.limit)
notify_queues(state)
staged
end)
rescue
error in [DBConnection.ConnectionError, Postgrex.Error] -> {:error, error}
end
defp stage_and_notify(_leader, state) do
if state.mode == :local, do: notify_queues(state)
{:ok, []}
end
defp notify_queues(%{conf: conf, mode: :global}) do
{:ok, queues} = Engine.check_available(conf)
payload = Enum.map(queues, &%{queue: &1})
Notifier.notify(conf, :insert, payload)
end
defp notify_queues(%{conf: conf, mode: :local}) do
match = [{{{conf.name, {:producer, :"$1"}}, :"$2", :_}, [], [{{:"$1", :"$2"}}]}]
for {queue, pid} <- Registry.select(match) do
send(pid, {:notification, :insert, %{"queue" => queue}})
end
:ok
end
# Helpers
defp schedule_staging(state) do
timer = Process.send_after(self(), :stage, state.interval)
%{state | timer: timer}
end
defp check_mode(state) do
next_mode =
case Notifier.status(state.conf) do
:clustered -> :global
:isolated -> :local
:solitary -> if Peer.leader?(state.conf), do: :global, else: :local
:unknown -> state.mode
end
if state.mode != next_mode do
:telemetry.execute([:oban, :stager, :switch], %{}, %{conf: state.conf, mode: next_mode})
end
%{state | mode: next_mode}
end
end