Packages
oban
2.19.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/queue/watchman.ex
defmodule Oban.Queue.Watchman do
@moduledoc false
use GenServer
alias Oban.Queue.Producer
alias __MODULE__, as: State
defstruct [:conf, :producer, :shutdown, interval: 10]
@spec child_spec(Keyword.t()) :: Supervisor.child_spec()
def child_spec(opts) do
shutdown = Keyword.fetch!(opts, :shutdown) + Keyword.get(opts, :interval, 10)
%{id: __MODULE__, start: {__MODULE__, :start_link, [opts]}, shutdown: shutdown}
end
@spec start_link(Keyword.t()) :: GenServer.on_start()
def start_link(opts) do
{name, opts} = Keyword.pop(opts, :name)
GenServer.start_link(__MODULE__, struct!(State, opts), name: name)
end
@impl GenServer
def init(state) do
Process.flag(:trap_exit, true)
{:ok, state}
end
@impl GenServer
def terminate(_reason, %State{} = state) do
# The producer may not exist, and we don't want to raise during shutdown.
:ok = Producer.shutdown(state.producer)
:ok = wait_for_executing(0, state)
catch
:exit, _reason -> :ok
end
defp wait_for_executing(elapsed, state) do
check = Producer.check(state.producer)
if check.running == [] or elapsed >= state.shutdown do
:telemetry.execute(
[:oban, :queue, :shutdown],
%{elapsed: elapsed, ellapsed: elapsed},
%{conf: state.conf, orphaned: check.running, queue: check.queue}
)
:ok
else
:ok = Process.sleep(state.interval)
wait_for_executing(elapsed + state.interval, state)
end
end
end