Packages
oban
1.2.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/executor.ex
defmodule Oban.Queue.Executor do
@moduledoc false
require Logger
alias Oban.{Breaker, Config, Job, Query, Worker}
@type success :: {:success, Job.t()}
@type failure :: {:failure, Job.t(), Worker.t(), atom(), term()}
@spec call(Job.t(), Config.t()) :: :success | :failure
def call(%Job{} = job, %Config{} = conf) do
start_time = System.system_time(:microsecond)
start_mono = System.monotonic_time(:microsecond)
:telemetry.execute([:oban, :started], %{start_time: start_time}, job_meta(job))
case safe_call(job) do
{:success, ^job} ->
Breaker.with_retry(fn -> report_success(conf, job, start_mono) end)
{:failure, ^job, kind, error, stack} ->
Breaker.with_retry(fn -> report_failure(conf, job, start_mono, kind, error, stack) end)
end
end
@spec safe_call(Job.t()) :: success() | failure()
def safe_call(%Job{} = job) do
case Job.worker_module(job) do
{:ok, worker} ->
case worker.timeout(job) do
:infinity ->
perform_inline(worker, job)
timeout when is_integer(timeout) ->
perform_timed(worker, job, timeout)
end
{:error, error, stacktrace} ->
{:failure, job, :error, error, stacktrace}
end
end
@spec perform_inline(Worker.t(), Job.t()) :: success() | failure()
def perform_inline(worker, %Job{args: args} = job) do
case worker.perform(args, job) do
:ok ->
{:success, job}
{:ok, _result} ->
{:success, job}
{:error, error} ->
{:failure, job, :error, error, current_stacktrace()}
returned ->
Logger.warn(fn ->
"""
Expected #{worker}.perform/2 to return :ok, `{:ok, value}`, or {:error, reason}. Instead
received:
#{inspect(returned, pretty: true)}
The job will be considered a success.
"""
end)
{:success, job}
end
rescue
exception ->
{:failure, job, :exception, exception, __STACKTRACE__}
catch
kind, value ->
{:failure, job, kind, value, __STACKTRACE__}
end
@spec perform_timed(Worker.t(), Job.t(), timeout()) :: success() | failure()
def perform_timed(worker, job, timeout) do
task = Task.async(fn -> perform_inline(worker, job) end)
case Task.yield(task, timeout) || Task.shutdown(task) do
{:ok, reply} ->
reply
nil ->
{:failure, job, :error, :timeout, current_stacktrace()}
end
end
@spec report_success(Config.t(), Job.t(), integer()) :: :success
def report_success(conf, job, start_mono) do
Query.complete_job(conf, job)
:telemetry.execute([:oban, :success], %{duration: duration(start_mono)}, job_meta(job))
:success
end
@spec report_failure(Config.t(), Job.t(), integer(), atom(), any(), list()) :: :failure
def report_failure(conf, job, start_mono, kind, error, stack) do
Query.retry_job(conf, job, worker_backoff(job), kind, error, stack)
meta =
job
|> job_meta()
|> Map.merge(%{kind: kind, error: error, stack: stack})
:telemetry.execute([:oban, :failure], %{duration: duration(start_mono)}, meta)
:failure
end
# Helpers
defp current_stacktrace do
self()
|> Process.info(:current_stacktrace)
|> elem(1)
end
defp duration(start_mono), do: System.monotonic_time(:microsecond) - start_mono
defp job_meta(job), do: Map.take(job, [:id, :args, :queue, :worker, :attempt, :max_attempts])
defp worker_backoff(%Job{attempt: attempt} = job) do
case Job.worker_module(job) do
{:ok, worker} -> worker.backoff(attempt)
_ -> Worker.default_backoff(attempt)
end
end
end