Packages
oban
1.0.0-rc.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/notifier.ex
defmodule Oban.Notifier do
@moduledoc false
# The notifier has several different responsibilities and some nuanced behavior:
#
# On Start:
# 1. Create a connection
# 2. Listen for insert/signal events
# 3. If connection fails then log the error, break the circuit, and attempt to connect later
#
# On Exit:
# 1. Trip the circuit breaker
# 2. Schedule a reconnect with backoff
#
# On Listen:
# 1. Put the producer into the listeners map
# 2. Monitor the pid so that we can clean up if the producer dies
#
# On Notification:
# 1. Iterate through the listeners and forward the message
use GenServer
import Oban.Breaker, only: [open_circuit: 1, trip_circuit: 2]
alias Oban.Config
alias Postgrex.Notifications
@type option :: {:name, module()} | {:conf, Config.t()}
@type channel :: :gossip | :insert | :signal
@type queue :: atom()
@mappings %{
gossip: "oban_gossip",
insert: "oban_insert",
signal: "oban_signal"
}
@channels Map.keys(@mappings)
defmodule State do
@moduledoc false
@enforce_keys [:conf]
defstruct [
:conf,
:conn,
:name,
circuit: :enabled,
listeners: %{}
]
end
defmacro gossip, do: @mappings[:gossip]
defmacro insert, do: @mappings[:insert]
defmacro signal, do: @mappings[:signal]
defguardp is_server(server) when is_pid(server) or is_atom(server)
@spec start_link([option]) :: GenServer.on_start()
def start_link(opts) do
name = Keyword.get(opts, :name, __MODULE__)
GenServer.start_link(__MODULE__, Map.new(opts), name: name)
end
@spec listen(module()) :: :ok
def listen(server, channels \\ @channels) when is_server(server) and is_list(channels) do
GenServer.call(server, {:listen, channels})
end
@impl GenServer
def init(opts) do
Process.flag(:trap_exit, true)
{:ok, struct!(State, opts), {:continue, :start}}
end
@impl GenServer
def handle_continue(:start, state) do
{:noreply, connect_and_listen(state)}
end
@impl GenServer
def handle_info({:DOWN, _ref, :process, pid, _reason}, %State{listeners: listeners} = state) do
{:noreply, %{state | listeners: Map.delete(listeners, pid)}}
end
def handle_info({:notification, _, _, prefixed_channel, payload}, state) do
[_prefix, channel] = String.split(prefixed_channel, ".")
decoded = Jason.decode!(payload)
for {pid, channels} <- state.listeners, channel in channels do
send(pid, {:notification, channel, decoded})
end
{:noreply, state}
end
def handle_info({:EXIT, _pid, error}, %State{} = state) do
state = trip_circuit(error, state)
{:noreply, %{state | conn: nil}}
end
def handle_info(:reset_circuit, %State{circuit: :disabled} = state) do
state =
state
|> open_circuit()
|> connect_and_listen()
{:noreply, state}
end
def handle_info(_message, state) do
{:noreply, state}
end
@impl GenServer
def handle_call({:listen, channels}, {pid, _}, %State{listeners: listeners} = state) do
if Map.has_key?(listeners, pid) do
{:reply, :ok, state}
else
Process.monitor(pid)
full_channels =
@mappings
|> Map.take(channels)
|> Map.values()
{:reply, :ok, %{state | listeners: Map.put(listeners, pid, full_channels)}}
end
end
defp connect_and_listen(%State{conf: conf, conn: nil} = state) do
case Notifications.start_link(conf.repo.config()) do
{:ok, conn} ->
Notifications.listen(conn, "#{conf.prefix}.#{gossip()}")
Notifications.listen(conn, "#{conf.prefix}.#{insert()}")
Notifications.listen(conn, "#{conf.prefix}.#{signal()}")
%{state | conn: conn}
{:error, error} ->
trip_circuit(error, state)
end
end
defp connect_and_listen(state), do: state
end