Packages
oban
2.4.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/midwife.ex
defmodule Oban.Midwife do
@moduledoc false
use GenServer
alias Oban.{Config, Notifier, Registry}
alias Oban.Queue.Supervisor, as: QueueSupervisor
@type option :: {:name, module()} | {:conf, Config.t()}
defmodule State do
@moduledoc false
defstruct [:conf]
end
@spec start_link([option]) :: GenServer.on_start()
def start_link(opts) do
{name, opts} = Keyword.pop(opts, :name, __MODULE__)
GenServer.start_link(__MODULE__, opts, name: name)
end
@impl GenServer
def init(opts) do
{:ok, struct!(State, opts), {:continue, :start}}
end
@impl GenServer
def handle_continue(:start, %State{conf: conf} = state) do
Notifier.listen(conf.name, [:signal])
{:noreply, state}
end
@impl GenServer
def handle_info({:notification, :signal, payload}, %State{conf: conf} = state) do
case payload do
%{"action" => "start", "queue" => queue, "limit" => limit} ->
Supervisor.start_child(Registry.via(conf.name), queue_spec(queue, limit, conf))
%{"action" => "stop", "queue" => queue} ->
%{id: child_id} = queue_spec(queue, 0, conf)
Supervisor.terminate_child(Registry.via(conf.name), child_id)
Supervisor.delete_child(Registry.via(conf.name), child_id)
_ ->
:ok
end
{:noreply, state}
end
def handle_info(_message, state) do
{:noreply, state}
end
defp queue_spec(queue, limit, conf) do
QueueSupervisor.child_spec({queue, limit: limit}, conf)
end
end