Packages
oban
2.17.8
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/sonar.ex
defmodule Oban.Sonar do
@moduledoc false
use GenServer
alias Oban.Notifier
alias __MODULE__, as: State
defstruct [
:conf,
:timer,
interval: :timer.seconds(15),
nodes: %{},
stale_mult: 2,
status: :unknown
]
@spec start_link(keyword()) :: GenServer.on_start()
def start_link(opts) do
{name, opts} = Keyword.pop(opts, :name)
conf = Keyword.fetch!(opts, :conf)
if conf.testing != :disabled do
:ignore
else
GenServer.start_link(__MODULE__, struct!(State, opts), name: name)
end
end
@impl GenServer
def init(state) do
Process.flag(:trap_exit, true)
{:ok, state, {:continue, :start}}
end
@impl GenServer
def terminate(_reason, state) do
if is_reference(state.timer), do: Process.cancel_timer(state.timer)
:ok
end
@impl GenServer
def handle_continue(:start, state) do
:ok = Notifier.listen(state.conf.name, :sonar)
:ok = Notifier.notify(state.conf, :sonar, %{node: state.conf.node, ping: :ping})
{:noreply, schedule_ping(state)}
end
@impl GenServer
def handle_call(:get_status, _from, state) do
{:reply, state.status, state}
end
def handle_call(:prune_nodes, _from, state) do
state = prune_stale_nodes(state)
{:reply, state.nodes, state}
end
@impl GenServer
def handle_info(:ping, state) do
:ok = Notifier.notify(state.conf, :sonar, %{node: state.conf.node, ping: true})
state =
state
|> prune_stale_nodes()
|> update_status()
|> schedule_ping()
{:noreply, state}
end
def handle_info({:notification, :sonar, %{"node" => node} = payload}, state) do
time = Map.get(payload, "time", System.system_time(:millisecond))
state =
state
|> Map.update!(:nodes, &Map.put(&1, node, time))
|> update_status()
{:noreply, state}
end
# Helpers
defp schedule_ping(state) do
timer = Process.send_after(self(), :ping, state.interval)
%{state | timer: timer}
end
defp update_status(state) do
node = state.conf.node
status =
case Map.keys(state.nodes) do
[] -> :isolated
[^node] -> :solitary
[_ | _] -> :clustered
end
if status != state.status do
:telemetry.execute([:oban, :notifier, :switch], %{}, %{conf: state.conf, status: status})
end
%{state | status: status}
end
defp prune_stale_nodes(state) do
stale = System.system_time(:millisecond) - state.interval * state.stale_mult
nodes = Map.reject(state.nodes, fn {_, recorded} -> recorded < stale end)
%{state | nodes: nodes}
end
end