Packages
oban
2.20.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/queue/producer.ex
defmodule Oban.Queue.Producer do
@moduledoc false
use GenServer
alias Oban.{Backoff, Engine, Notifier, PerformError, TimeoutError, Worker}
alias Oban.Queue.Executor
alias __MODULE__, as: State
require Logger
defstruct [
:conf,
:foreman,
:meta,
:name,
:dispatch_timer,
:refresh_timer,
dispatch_cooldown: 5,
running: %{}
]
@spec start_link([Keyword.t()]) :: GenServer.on_start()
def start_link(opts) do
GenServer.start_link(__MODULE__, opts, name: opts[:name])
end
@spec check(GenServer.server()) :: Oban.queue_state()
def check(producer) do
GenServer.call(producer, :check)
end
@spec shutdown(GenServer.server()) :: :ok
def shutdown(producer) do
GenServer.call(producer, :shutdown)
end
# Callbacks
@impl GenServer
def init(opts) do
Process.flag(:trap_exit, true)
{base_opts, meta_opts} = Keyword.split(opts, [:conf, :foreman, :name, :dispatch_cooldown])
state = struct!(State, base_opts)
{:ok, state, {:continue, {:start, meta_opts}}}
end
@impl GenServer
def terminate(_reason, %State{} = state) do
if is_reference(state.refresh_timer), do: Process.cancel_timer(state.refresh_timer)
if is_reference(state.dispatch_timer), do: Process.cancel_timer(state.dispatch_timer)
:ok
end
@impl GenServer
def handle_continue({:start, meta_opts}, %State{} = state) do
{:ok, meta} = Engine.init(state.conf, meta_opts)
:ok = Notifier.listen(state.conf.name, [:insert, :signal])
{:noreply, schedule_refresh(%{state | meta: meta})}
end
@impl GenServer
def handle_info({ref, _val}, %State{} = state) when is_reference(ref) do
state =
state
|> release_ref(ref)
|> schedule_dispatch()
{:noreply, state}
end
def handle_info({:DOWN, ref, :process, pid, reason}, %State{} = state) do
{^pid, exec} = Map.get(state.running, ref)
# The worker module is resolved after storing the exec struct. We try to resolve it again in
# order to get a worker's custom `backoff/1`.
exec =
case Worker.from_string(exec.job.worker) do
{:ok, worker} -> %{exec | worker: worker}
{:error, _error} -> exec
end
exec =
case reason do
%TimeoutError{} ->
%{exec | kind: :error, error: reason, state: :failure}
:shutdown ->
result = {:cancel, :shutdown}
reason = PerformError.exception({exec.job.worker, result})
%{exec | result: result, error: reason, state: :cancelled}
{error, stack} ->
%{exec | kind: {:EXIT, pid}, error: error, stacktrace: stack, state: :failure}
_ ->
%{exec | kind: {:EXIT, pid}, error: reason, state: :failure}
end
Task.Supervisor.async_nolink(state.foreman, fn ->
exec =
exec
|> Executor.normalize_state()
|> Executor.record_finished()
|> Executor.cancel_timeout()
Backoff.with_retry(fn -> Executor.report_finished(exec) end)
end)
state =
state
|> release_ref(ref)
|> schedule_dispatch()
{:noreply, state}
end
def handle_info({:notification, :insert, %{"queue" => queue}}, %State{} = state) do
state = if state.meta.queue == queue, do: schedule_dispatch(state), else: state
{:noreply, state}
end
def handle_info({:notification, :signal, payload}, %State{} = state) do
queue = state.meta.queue
meta =
case payload do
%{"action" => "pause", "queue" => "*"} ->
Engine.put_meta(state.conf, state.meta, :paused, true)
%{"action" => "resume", "queue" => "*"} ->
Engine.put_meta(state.conf, state.meta, :paused, false)
%{"action" => "pause", "queue" => ^queue} ->
Engine.put_meta(state.conf, state.meta, :paused, true)
%{"action" => "resume", "queue" => ^queue} ->
Engine.put_meta(state.conf, state.meta, :paused, false)
%{"action" => "scale", "queue" => ^queue} ->
payload
|> Map.drop(~w(action ident node queue))
|> Enum.reduce(state.meta, fn {key, val}, meta ->
Engine.put_meta(state.conf, meta, String.to_existing_atom(key), val)
end)
%{"action" => "pkill", "job_id" => jid} ->
for {ref, {pid, exec}} <- state.running, exec.job.id == jid do
pkill(ref, pid, state)
end
state.meta
_ ->
state.meta
end
{:noreply, schedule_dispatch(%{state | meta: meta})}
end
def handle_info(:dispatch, %State{} = state) do
{:noreply, dispatch(%{state | dispatch_timer: nil})}
end
def handle_info(:refresh, %State{} = state) do
meta = Engine.refresh(state.conf, state.meta)
{:noreply, schedule_refresh(%{state | meta: meta})}
end
def handle_info(message, state) do
Logger.warning(
message: "Received unexpected message: #{inspect(message)}",
source: :oban,
module: __MODULE__
)
{:noreply, state}
end
@impl GenServer
def handle_call(:check, _from, %State{} = state) do
meta = Engine.check_meta(state.conf, state.meta, state.running)
{:reply, meta, state}
end
def handle_call({:put_meta, key, value}, _from, %State{} = state) do
meta = Engine.put_meta(state.conf, state.meta, key, value)
{:reply, meta, %{state | meta: meta}}
end
def handle_call(:shutdown, _from, %State{} = state) do
meta = Engine.shutdown(state.conf, state.meta)
{:reply, :ok, %{state | meta: meta}}
end
# Killing
defp pkill(ref, pid, %State{} = state) do
case Task.Supervisor.terminate_child(state.foreman, pid) do
:ok ->
state
{:error, :not_found} ->
release_ref(state, ref)
end
end
# Dispatching
defp schedule_dispatch(%State{} = state) do
if is_reference(state.dispatch_timer) do
state
else
timer = Process.send_after(self(), :dispatch, state.dispatch_cooldown)
%{state | dispatch_timer: timer}
end
end
defp schedule_refresh(%State{} = state) do
cooldown = Backoff.jitter(state.meta.refresh_interval, mode: :dec, mult: 0.5)
timer = Process.send_after(self(), :refresh, cooldown)
%{state | refresh_timer: timer}
end
defp dispatch(%State{} = state) do
tele_meta = %{conf: state.conf, queue: state.meta.queue}
{meta, dispatched} =
:telemetry.span([:oban, :producer], tele_meta, fn ->
{meta, dispatched} = start_jobs(state)
{{meta, dispatched}, Map.put(tele_meta, :dispatched_count, map_size(dispatched))}
end)
%{state | meta: meta, running: Map.merge(state.running, dispatched)}
end
defp start_jobs(%State{} = state) do
{:ok, {meta, jobs}} = Engine.fetch_jobs(state.conf, state.meta, state.running)
dispatched =
for job <- jobs, into: %{} do
exec = Executor.new(state.conf, job)
task = Task.Supervisor.async_nolink(state.foreman, Executor, :call, [exec])
{task.ref, {task.pid, exec}}
end
{meta, dispatched}
end
defp release_ref(%State{} = state, ref) do
Process.demonitor(ref, [:flush])
%{state | running: Map.delete(state.running, ref)}
end
end