Packages
oban
0.8.0
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
use GenServer
alias Oban.{Config, Query}
alias Postgrex.Notifications
@type option :: {:name, module()} | {:conf, Config.t()}
@type channel :: :insert | :signal | :update
@type queue :: atom()
@mappings %{
insert: "oban_insert",
signal: "oban_signal",
update: "oban_update"
}
@channels Map.keys(@mappings)
defmodule State do
@moduledoc false
@enforce_keys [:conf]
defstruct [:conf]
end
defmacro insert, do: @mappings[:insert]
defmacro signal, do: @mappings[:signal]
defmacro update, do: @mappings[:update]
@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(), binary(), channel()) :: :ok
def listen(server, prefix, channel) when is_binary(prefix) and channel in @channels do
server
|> conn_name()
|> Notifications.listen("#{prefix}.#{@mappings[channel]}")
:ok
end
@spec notify(module(), channel(), term()) :: :ok
def notify(server, channel, payload) when channel in @channels do
GenServer.call(server, {:notify, @mappings[channel], payload})
end
@spec pause_queue(module(), queue()) :: :ok
def pause_queue(server, queue) when is_atom(queue) do
notify(server, :signal, %{action: :pause, queue: queue})
end
@spec resume_queue(module(), queue()) :: :ok
def resume_queue(server, queue) when is_atom(queue) do
notify(server, :signal, %{action: :resume, queue: queue})
end
@spec scale_queue(module(), queue(), pos_integer()) :: :ok
def scale_queue(server, queue, scale)
when is_atom(queue) and is_integer(scale) and scale > 0 do
notify(server, :signal, %{action: :scale, queue: queue, scale: scale})
end
@spec kill_job(module(), pos_integer()) :: :ok
def kill_job(server, job_id) when is_integer(job_id) do
notify(server, :signal, %{action: :pkill, job_id: job_id})
end
@impl GenServer
def init(%{conf: conf, name: name}) do
{:ok, %State{conf: conf}, {:continue, {:start, name}}}
end
@impl GenServer
def handle_continue({:start, name}, %State{conf: conf} = state) do
conn_conf = Keyword.put(conf.repo.config(), :name, conn_name(name))
{:ok, _} = Notifications.start_link(conn_conf)
{:noreply, state}
end
@impl GenServer
def handle_call({:notify, channel, payload}, _from, %State{conf: conf} = state) do
# Unlike `listen` the schema namespacing takes place in the Query module.
:ok = Query.notify(conf, channel, Jason.encode!(payload))
{:reply, :ok, state}
end
defp conn_name(name), do: Module.concat(name, "Conn")
end