Packages
oban
2.3.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/queue/drainer.ex
defmodule Oban.Queue.Drainer do
@moduledoc false
alias Oban.{Config, Query}
alias Oban.Queue.Executor
@unlimited 100_000_000
@far_future DateTime.from_unix!(9_999_999_999)
@type drain_option ::
{:queue, binary() | atom()}
| {:with_scheduled, boolean()}
| {:with_safety, boolean()}
@type drain_result :: %{success: non_neg_integer(), failure: non_neg_integer()}
@spec drain(Config.t(), [drain_option()]) :: drain_result()
def drain(%Config{} = conf, [_ | _] = opts) do
queue =
opts
|> Keyword.fetch!(:queue)
|> to_string()
if Keyword.get(opts, :with_scheduled, false), do: schedule_jobs(conf, queue)
conf
|> fetch_jobs(queue)
|> Enum.reduce(%{failure: 0, success: 0}, fn job, acc ->
result =
conf
|> Executor.new(job)
|> Executor.put(:safe, Keyword.get(opts, :with_safety, true))
|> Executor.call()
Map.update(acc, result, 1, &(&1 + 1))
end)
end
defp schedule_jobs(conf, queue) do
Query.stage_scheduled_jobs(conf, queue, max_scheduled_at: @far_future)
end
defp fetch_jobs(conf, queue) do
{:ok, jobs} = Query.fetch_available_jobs(conf, queue, "draining", @unlimited)
jobs
end
end