Packages
oban
1.0.0-rc.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/queue/executor.ex
defmodule Oban.Queue.Executor do
@moduledoc false
alias Oban.{Config, Job, Query, Worker}
@spec child_spec(Job.t(), Config.t()) :: Supervisor.child_spec()
def child_spec(job, conf) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [job, conf]},
type: :worker,
restart: :temporary
}
end
@spec start_link(Job.t(), Config.t()) :: {:ok, pid()}
def start_link(%Job{} = job, %Config{} = conf) do
Task.start_link(__MODULE__, :call, [job, conf])
end
@spec call(Job.t(), Config.t()) :: :success | :failure
def call(%Job{} = job, %Config{} = conf) do
{duration, return} = :timer.tc(__MODULE__, :safe_call, [job])
case return do
{:success, ^job} ->
Query.complete_job(conf, job)
report(:success, duration, job, %{})
{:failure, ^job, kind, error, stack} ->
Query.retry_job(conf, job, worker_backoff(job), format_blamed(kind, error, stack))
report(:failure, duration, job, %{kind: kind, error: error, stack: stack})
end
end
@doc false
def safe_call(%Job{args: args, worker: worker} = job) do
worker
|> to_module()
|> apply(:perform, [args, job])
|> case do
{:error, error} ->
{:current_stacktrace, stacktrace} = Process.info(self(), :current_stacktrace)
{:failure, job, :error, error, stacktrace}
_ ->
{:success, job}
end
rescue
exception ->
{:failure, job, :exception, exception, __STACKTRACE__}
catch
kind, value ->
{:failure, job, kind, value, __STACKTRACE__}
end
# Helpers
defp to_module(worker) when is_binary(worker) do
worker
|> String.split(".")
|> Module.safe_concat()
end
defp to_module(worker) when is_atom(worker), do: worker
# While it is slightly wasteful, we have to convert the worker to a module again outside of
# `safe_call/1`. There is a possibility that the worker module can't be found at all and we
# need to fall back to a default implementation.
defp worker_backoff(%Job{attempt: attempt, worker: worker}) do
module =
try do
to_module(worker)
rescue
ArgumentError -> nil
end
if function_exported?(module, :backoff, 1) do
module.backoff(attempt)
else
Worker.default_backoff(attempt)
end
end
defp format_blamed(:exception, error, stack), do: format_blamed(:error, error, stack)
defp format_blamed(kind, error, stack) do
{blamed, stack} = Exception.blame(kind, error, stack)
Exception.format(kind, blamed, stack)
end
defp report(event, duration, job, meta) do
meta =
job
|> Map.take([:id, :args, :queue, :worker, :attempt, :max_attempts])
|> Map.merge(meta)
:telemetry.execute([:oban, event], %{duration: duration}, meta)
event
end
end