Packages
oban
2.15.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.
## Migrating from `Oban.Notifiers.Postgres`
After switching from `Oban.Notifiers.Postgres`, you may remove the unused `oban_notify` trigger.
Use the following migration to drop the trigger while retaining the `oban_jobs_notify` function:
```elixir
defmodule MyApp.Repo.Migrations.DropObanJobsNotifyTrigger do
use Ecto.Migration
def change do
execute(
"DROP TRIGGER IF EXISTS oban_notify ON public.oban_jobs",
"CREATE TRIGGER oban_notify AFTER INSERT ON public.oban_jobs FOR EACH ROW EXECUTE PROCEDURE public.oban_jobs_notify()"
)
end
end
```
[de]: https://elixir-lang.org/getting-started/mix-otp/distributed-tasks.html#our-first-distributed-code
"""
@behaviour Oban.Notifier
use GenServer
alias Oban.Notifier
defmodule State do
@moduledoc false
defstruct [:conf, :name, listeners: %{}]
end
@impl Notifier
def start_link(opts) do
{name, opts} = Keyword.pop(opts, :name, __MODULE__)
GenServer.start_link(__MODULE__, opts, name: name)
end
@impl Notifier
def listen(server, channels) do
GenServer.call(server, {:listen, channels})
end
@impl Notifier
def unlisten(server, channels) do
GenServer.call(server, {:unlisten, channels})
end
@impl 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
listeners = for {pid, channels} <- state.listeners, channel in channels, do: pid
Notifier.relay(state.conf, listeners, channel, payload)
{: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
end