Packages
oban
2.17.2
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 (Process Groups)][pg] based notifier implementation that runs with Distributed Erlang.
> #### Distributed Erlang {: .info}
>
> PG requires a functional [Distributed Erlang][de] cluster to broadcast between nodes. If your
> application isn't clustered, then you should consider an alternative notifier such as
> `Oban.Notifiers.Postgres`
## Usage
Specify the `PG` notifier in your Oban configuration:
```elixir
config :my_app, Oban,
notifier: Oban.Notifiers.PG,
...
```
[pg]: https://www.erlang.org/doc/man/pg.html
[de]: https://elixir-lang.org/getting-started/mix-otp/distributed-tasks.html#our-first-distributed-code
"""
@behaviour Oban.Notifier
use GenServer
alias Oban.Notifier
defstruct [:conf, listeners: %{}]
@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 %{conf: conf} <- get_state(server) do
pids = :pg.get_members(__MODULE__, conf.prefix)
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!(__MODULE__, opts)
put_state(state)
:pg.start_link(__MODULE__)
:pg.join(__MODULE__, state.conf.prefix, self())
{:ok, state}
end
defp put_state(state) do
Registry.update_value(Oban.Registry, {state.conf.name, Oban.Notifier}, fn _ -> state end)
end
defp get_state(server) do
[name] = Registry.keys(Oban.Registry, server)
case Oban.Registry.lookup(name) do
{_pid, state} -> state
nil -> :error
end
end
@impl GenServer
def handle_call({:listen, channels}, {pid, _}, %{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, _}, %{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
@impl GenServer
def handle_info({:notification, channel, payload}, 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({:DOWN, _ref, :process, pid, _reason}, state) do
{:noreply, %{state | listeners: Map.delete(state.listeners, pid)}}
end
def handle_info(_message, state) do
{:noreply, state}
end
## Message Helpers
defp payload_to_messages(channel, payload) do
Enum.map(payload, &{:notification, channel, &1})
end
end