Packages
oban
2.14.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/notifiers/pg.ex
defmodule Oban.Notifiers.PG do
@moduledoc """
A PG/PG2 based notifier implementation that runs with Distributed Erlang.
Out of the box, Oban uses PostgreSQL's `LISTEN/NOTIFY` for PubSub. For most applications, that
is fine, but Postgres-based PubSub isn't sufficient in some circumstances. In particular,
Postgres notifications won't work when your application connects through PGbouncer in
_transaction_ or _statement_ mode.
_Note: You must be using [Distributed Erlang][de] to use the PG notifier._
## Usage
Specify the `PG` notifier in your Oban configuration:
```elixir
config :my_app, Oban,
notifier: Oban.Notifiers.PG,
...
```
## Implementation Notes
* The notifier will use `pg` if available (OTP 23+) or fall back to `pg2` for
older OTP releases.
* Like the Postgres implementation, notifications are namespaced by `prefix`.
* For compatibility, message payloads are always serialized to JSON before
broadcast and deserialized before relay to local processes.
[de]: https://elixir-lang.org/getting-started/mix-otp/distributed-tasks.html#our-first-distributed-code
"""
@behaviour Oban.Notifier
use GenServer
alias Oban.Config
defmodule State do
@moduledoc false
defstruct [:conf, :name, listeners: %{}]
end
@impl Oban.Notifier
def start_link(opts) do
{name, opts} = Keyword.pop(opts, :name, __MODULE__)
GenServer.start_link(__MODULE__, opts, name: name)
end
@impl Oban.Notifier
def listen(server, channels) do
GenServer.call(server, {:listen, channels})
end
@impl Oban.Notifier
def unlisten(server, channels) do
GenServer.call(server, {:unlisten, channels})
end
@impl Oban.Notifier
def notify(server, channel, payload) do
with %State{} = state <- GenServer.call(server, :get_state),
[_ | _] = pids <- members(state.conf.prefix) do
for pid <- pids, message <- payload_to_messages(channel, payload) do
send(pid, message)
end
:ok
end
end
@impl GenServer
def init(opts) do
state = struct!(State, opts)
start_pg()
:ok = join(state.conf.prefix)
{:ok, 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)
{:reply, :ok, %{state | listeners: Map.put(listeners, pid, channels)}}
end
end
def handle_call({:unlisten, channels}, {pid, _}, %State{listeners: listeners} = state) do
orig_channels = Map.get(listeners, pid, [])
listeners =
case orig_channels -- channels do
[] -> Map.delete(listeners, pid)
new_channels -> Map.put(listeners, pid, new_channels)
end
{:reply, :ok, %{state | listeners: listeners}}
end
def handle_call(:get_state, _from, state), do: {:reply, state, state}
@impl GenServer
def handle_info({:notification, channel, payload}, %State{} = state) do
decoded = Jason.decode!(payload)
if in_scope?(decoded, state.conf) do
for {pid, channels} <- state.listeners, channel in channels do
send(pid, {:notification, channel, decoded})
end
end
{:noreply, state}
end
def handle_info(_message, state) do
{:noreply, state}
end
## PG Helpers
if Code.ensure_loaded?(:pg) do
defp start_pg do
:pg.start_link(__MODULE__)
end
defp members(prefix) do
:pg.get_members(__MODULE__, prefix)
end
defp join(prefix) do
:ok = :pg.join(__MODULE__, prefix, self())
end
else
defp start_pg, do: :ok
defp members(prefix) do
:pg2.get_members(namespace(prefix))
end
defp join(prefix) do
namespace = namespace(prefix)
:ok = :pg2.create(namespace)
:ok = :pg2.join(namespace, self())
end
defp namespace(prefix), do: {:oban, prefix}
end
## Message Helpers
defp payload_to_messages(channel, payload) do
Enum.map(payload, &{:notification, channel, &1})
end
defp in_scope?(%{"ident" => "any"}, _conf), do: true
defp in_scope?(%{"ident" => ident}, conf), do: Config.match_ident?(conf, ident)
defp in_scope?(_payload, _conf), do: true
end