Packages
oban
2.10.1
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/breaker.ex
defmodule Oban.Breaker do
@moduledoc false
alias Oban.{Config, Telemetry}
require Logger
@type state_struct :: %{
:circuit => :enabled | :disabled,
:conf => Config.t(),
:name => GenServer.name(),
:reset_timer => reference(),
optional(atom()) => any()
}
@max_retries 10
@min_delay 100
@default_jitter_mult 0.10
@default_jitter_mode :both
defmacro trip_errors, do: [DBConnection.ConnectionError, Postgrex.Error]
@spec trip_circuit(Exception.t(), Exception.stacktrace(), state_struct()) :: state_struct()
def trip_circuit(exception, stacktrace, state) do
meta = %{
config: state.conf,
error: exception,
kind: :error,
reason: exception,
message: error_message(exception),
name: state.name,
stacktrace: stacktrace
}
Telemetry.execute([:oban, :circuit, :trip], %{}, meta)
if is_reference(state.reset_timer), do: Process.cancel_timer(state.reset_timer)
reset_timer = Process.send_after(self(), :reset_circuit, state.conf.circuit_backoff)
%{state | circuit: :disabled, reset_timer: reset_timer}
end
@spec open_circuit(state_struct()) :: state_struct()
def open_circuit(%{circuit: _, name: name, conf: conf} = state) do
Telemetry.execute([:oban, :circuit, :open], %{}, %{name: name, conf: conf})
%{state | circuit: :enabled}
end
@spec jitter(time :: pos_integer(), opts :: Keyword.t()) :: pos_integer()
def jitter(time, opts \\ []) do
mode = Keyword.get(opts, :mode, @default_jitter_mode)
mult = Keyword.get(opts, :mult, @default_jitter_mult)
diff = trunc(:rand.uniform() * mult * time)
case mode do
:inc ->
time + diff
:dec ->
time - diff
:both ->
if :rand.uniform() >= 0.5 do
time + diff
else
time - diff
end
end
end
@spec with_retry(fun(), integer()) :: term()
def with_retry(fun, retries \\ 0)
def with_retry(fun, @max_retries), do: fun.()
def with_retry(fun, retries) do
fun.()
catch
_kind, _value -> lazy_retry(fun, retries)
end
defp lazy_retry(fun, retries) do
time = @min_delay * :math.pow(2, retries)
time
|> trunc()
|> jitter()
|> Process.sleep()
with_retry(fun, retries + 1)
end
defp error_message(%Postgrex.Error{} = exception) do
Postgrex.Error.message(exception)
end
defp error_message(%DBConnection.ConnectionError{} = exception) do
DBConnection.ConnectionError.message(exception)
end
defp error_message(exception), do: inspect(exception)
end