Packages
oban
2.4.3
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.{
Breaker,
Config,
CrashError,
Job,
PerformError,
Query,
Telemetry,
TimeoutError,
Worker
}
require Logger
@type success :: {:success, Job.t()}
@type failure :: {:failure, Job.t(), Worker.t(), atom(), term()}
@type t :: %__MODULE__{
conf: Config.t(),
duration: pos_integer(),
job: Job.t(),
kind: any(),
meta: map(),
queue_time: integer(),
start_mono: integer(),
start_time: integer(),
stop_mono: integer(),
worker: Worker.t(),
safe: boolean(),
snooze: pos_integer(),
stacktrace: Exception.stacktrace(),
state: :unset | :discard | :failure | :success | :snoozed
}
@enforce_keys [:conf, :job]
defstruct [
:conf,
:error,
:job,
:meta,
:snooze,
:start_mono,
:start_time,
:stop_mono,
:worker,
safe: true,
duration: 0,
kind: :error,
queue_time: 0,
stacktrace: [],
state: :unset
]
@spec new(Config.t(), Job.t()) :: t()
def new(%Config{} = conf, %Job{} = job) do
struct!(__MODULE__,
conf: conf,
job: job,
meta: event_metadata(conf, job),
start_mono: System.monotonic_time(),
start_time: System.system_time()
)
end
@spec put(t(), :safe, boolean()) :: t()
def put(%__MODULE__{} = exec, :safe, value) when is_boolean(value) do
%{exec | safe: value}
end
@spec call(t()) :: :success | :failure
def call(%__MODULE__{} = exec) do
exec =
exec
|> record_started()
|> resolve_worker()
|> perform()
|> record_finished()
Breaker.with_retry(fn -> report_finished(exec).state end)
end
def record_started(%__MODULE__{} = exec) do
Telemetry.execute([:oban, :job, :start], %{system_time: exec.start_time}, exec.meta)
exec
end
@spec resolve_worker(t()) :: t()
def resolve_worker(%__MODULE__{} = exec) do
case Worker.from_string(exec.job.worker) do
{:ok, worker} ->
%{exec | worker: worker}
{:error, error} ->
unless exec.safe, do: raise(error)
%{exec | state: :failure, error: error}
end
end
@spec perform(t()) :: t()
def perform(%__MODULE__{state: :unset} = exec) do
case exec.worker.timeout(exec.job) do
:infinity ->
perform_inline(exec)
timeout when is_integer(timeout) ->
perform_timed(exec, timeout)
end
end
def perform(%__MODULE__{} = exec), do: exec
@spec record_finished(t()) :: t()
def record_finished(%__MODULE__{} = exec) do
stop_mono = System.monotonic_time()
duration = stop_mono - exec.start_mono
queue_time = DateTime.diff(exec.job.attempted_at, exec.job.scheduled_at, :nanosecond)
%{exec | duration: duration, queue_time: queue_time, stop_mono: stop_mono}
end
@spec report_finished(t()) :: t()
def report_finished(%__MODULE__{} = exec) do
case exec.state do
:success ->
Query.complete_job(exec.conf, exec.job)
execute_stop(exec)
:failure ->
job = job_with_unsaved_error(exec)
Query.retry_job(exec.conf, job, backoff(exec.worker, job))
execute_exception(exec)
:snoozed ->
Query.snooze_job(exec.conf, exec.job, exec.snooze)
execute_stop(exec)
:discard ->
Query.discard_job(exec.conf, job_with_unsaved_error(exec))
execute_stop(exec)
end
exec
end
defp perform_inline(%{safe: true} = exec) do
perform_inline(%{exec | safe: false})
rescue
error ->
%{exec | state: :failure, error: error, stacktrace: __STACKTRACE__}
catch
kind, reason ->
error = CrashError.exception({kind, reason, __STACKTRACE__})
%{exec | state: :failure, error: error, stacktrace: __STACKTRACE__}
end
defp perform_inline(%{worker: worker, job: job} = exec) do
case worker.perform(job) do
:ok ->
%{exec | state: :success}
{:ok, _result} ->
%{exec | state: :success}
:discard ->
%{exec | state: :discard, error: PerformError.exception({worker, :discard})}
{:discard, reason} ->
%{exec | state: :discard, error: PerformError.exception({worker, {:discard, reason}})}
{:error, reason} ->
%{exec | state: :failure, error: PerformError.exception({worker, {:error, reason}})}
{:snooze, seconds} ->
%{exec | state: :snoozed, snooze: seconds}
returned ->
Logger.warn(fn ->
"""
Expected #{worker}.perform/1 to return:
- `:ok`
- `:discard`
- `{:ok, value}`
- `{:error, reason}`,
- `{:discard, reason}`
- `{:snooze, seconds}`
Instead received:
#{inspect(returned, pretty: true)}
The job will be considered a success.
"""
end)
%{exec | state: :success}
end
end
defp perform_timed(exec, timeout) do
task = Task.async(fn -> perform_inline(exec) end)
case Task.yield(task, timeout) || Task.shutdown(task) do
{:ok, reply} ->
reply
nil ->
error = TimeoutError.exception({exec.worker, timeout})
%{exec | state: :failure, error: error}
end
end
defp backoff(nil, job), do: Worker.backoff(job)
defp backoff(worker, job), do: worker.backoff(job)
defp execute_stop(exec) do
measurements = %{duration: exec.duration, queue_time: exec.queue_time}
metadata = Map.put(exec.meta, :state, exec.state)
Telemetry.execute([:oban, :job, :stop], measurements, metadata)
end
defp execute_exception(exec) do
measurements = %{duration: exec.duration, queue_time: exec.queue_time}
meta =
Map.merge(exec.meta, %{
kind: exec.kind,
error: exec.error,
stacktrace: exec.stacktrace,
state: exec.state
})
Telemetry.execute([:oban, :job, :exception], measurements, meta)
end
defp event_metadata(conf, job) do
job
|> Map.take([:id, :args, :queue, :worker, :attempt, :max_attempts, :tags])
|> Map.put(:job, job)
|> Map.put(:prefix, conf.prefix)
|> Map.put(:conf, conf)
end
defp job_with_unsaved_error(%__MODULE__{} = exec) do
%{
exec.job
| unsaved_error: %{kind: exec.kind, reason: exec.error, stacktrace: exec.stacktrace}
}
end
end